diff --git a/tests/go.mod b/tests/go.mod index 358db52..d71fc81 100644 --- a/tests/go.mod +++ b/tests/go.mod @@ -5,24 +5,23 @@ go 1.26 toolchain go1.26.4 require ( - connectrpc.com/connect v1.20.0 github.com/google/uuid v1.6.0 - github.com/roadrunner-server/api-go/v6 v6.0.0-beta.12 + github.com/roadrunner-server/api-go/v6 v6.0.0-beta.12.0.20260714200341-93604e5012d4 github.com/roadrunner-server/api-plugins/v6 v6.0.0-beta.2 github.com/roadrunner-server/config/v6 v6.0.0-beta.3 github.com/roadrunner-server/endure/v2 v2.6.2 - github.com/roadrunner-server/http/v6 v6.0.0-beta.7 - github.com/roadrunner-server/informer/v6 v6.0.0-beta.2 - github.com/roadrunner-server/jobs/v6 v6.0.0-beta.7 - github.com/roadrunner-server/kv/v6 v6.0.0-beta.5 + github.com/roadrunner-server/goridge/v4 v4.0.0-beta.2.0.20260714195909-75e9ece43063 + github.com/roadrunner-server/http/v6 v6.0.0-beta.7.0.20260714202301-3d2c74eb5e61 + github.com/roadrunner-server/informer/v6 v6.0.0-beta.2.0.20260714201850-2854c943433a + github.com/roadrunner-server/jobs/v6 v6.0.0-beta.7.0.20260714202317-a904df360b84 + github.com/roadrunner-server/kv/v6 v6.0.0-beta.5.0.20260714202229-7ef2a556a8ad github.com/roadrunner-server/logger/v6 v6.0.0-beta.3 github.com/roadrunner-server/memory/v6 v6.0.0-beta.4 - github.com/roadrunner-server/resetter/v6 v6.0.0-beta.3 - github.com/roadrunner-server/rpc/v6 v6.0.0-beta.4 + github.com/roadrunner-server/resetter/v6 v6.0.0-beta.3.0.20260714201533-91a174bb65d4 + github.com/roadrunner-server/rpc/v6 v6.0.0-beta.4.0.20260714200548-15b82bc47898 github.com/roadrunner-server/server/v6 v6.0.0-beta.6 github.com/stretchr/testify v1.11.1 go.opentelemetry.io/otel/sdk v1.44.0 - golang.org/x/net v0.56.0 google.golang.org/genproto v0.0.0-20260630182238-925bb5da69e7 google.golang.org/protobuf v1.36.11 ) @@ -30,7 +29,6 @@ require ( replace github.com/roadrunner-server/memory/v6 => ../ require ( - connectrpc.com/grpcreflect v1.3.0 // indirect github.com/beorn7/perks v1.0.1 // indirect github.com/caddyserver/certmagic v0.25.4 // indirect github.com/caddyserver/zerossl v0.1.5 // indirect @@ -63,7 +61,6 @@ require ( github.com/roadrunner-server/context v1.3.0 // indirect github.com/roadrunner-server/errors v1.5.0 // indirect github.com/roadrunner-server/events v1.0.1 // indirect - github.com/roadrunner-server/goridge/v4 v4.0.0-beta.2 // indirect github.com/roadrunner-server/pool/v2 v2.0.0-beta.1 // indirect github.com/roadrunner-server/priority_queue v1.0.6 // indirect github.com/roadrunner-server/tcplisten v1.5.2 // indirect @@ -91,11 +88,10 @@ require ( go.yaml.in/yaml/v3 v3.0.4 // indirect golang.org/x/crypto v0.53.0 // indirect golang.org/x/mod v0.37.0 // indirect + golang.org/x/net v0.56.0 // indirect golang.org/x/sync v0.21.0 // indirect golang.org/x/sys v0.46.0 // indirect golang.org/x/text v0.38.0 // indirect golang.org/x/tools v0.47.0 // indirect - google.golang.org/genproto/googleapis/rpc v0.0.0-20260630182238-925bb5da69e7 // indirect - google.golang.org/grpc v1.82.0 // indirect gopkg.in/yaml.v3 v3.0.1 // indirect ) diff --git a/tests/go.sum b/tests/go.sum index 8de9219..c156c58 100644 --- a/tests/go.sum +++ b/tests/go.sum @@ -1,9 +1,5 @@ code.pfad.fr/check v1.1.0 h1:GWvjdzhSEgHvEHe2uJujDcpmZoySKuHQNrZMfzfO0bE= code.pfad.fr/check v1.1.0/go.mod h1:NiUH13DtYsb7xp5wll0U4SXx7KhXQVCtRgdC96IPfoM= -connectrpc.com/connect v1.20.0 h1:6TNDAB+WeNd2uolWNlYczB5E0KNNaVMNUEx8JEUsPmQ= -connectrpc.com/connect v1.20.0/go.mod h1:A2ygJrukXwWy32vkCAAHNVguZrqZ+jeZ9rGRnGR4dN4= -connectrpc.com/grpcreflect v1.3.0 h1:Y4V+ACf8/vOb1XOc251Qun7jMB75gCUNw6llvB9csXc= -connectrpc.com/grpcreflect v1.3.0/go.mod h1:nfloOtCS8VUQOQ1+GTdFzVg2CJo4ZGaat8JIovCtDYs= github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM= github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw= github.com/caddyserver/certmagic v0.25.4 h1:8eIXh0HC3MsGnNo8One+BCxMGTbe5zb/oz+2KsxBFQg= @@ -34,8 +30,6 @@ github.com/go-ole/go-ole v1.3.0 h1:Dt6ye7+vXGIKZ7Xtk4s6/xVdGDQynvom7xCFEdWr6uE= github.com/go-ole/go-ole v1.3.0/go.mod h1:5LS6F96DhAwUc7C+1HLexzMXY1xGRSryjyPPKW6zv78= github.com/go-viper/mapstructure/v2 v2.5.0 h1:vM5IJoUAy3d7zRSVtIwQgBj7BiWtMPfmPEgAXnvj1Ro= github.com/go-viper/mapstructure/v2 v2.5.0/go.mod h1:oJDH3BJKyqBA2TXFhDsKDGDTlndYOZ6rGS0BRZIxGhM= -github.com/golang/protobuf v1.5.4 h1:i7eJL8qZTpSEXOPTxNKhASYpMn+8e5Q6AdndVa1dWek= -github.com/golang/protobuf v1.5.4/go.mod h1:lnTiLA8Wa4RWRcIUkrtSVa5nRhsEGBg48fD6rSs7xps= github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8= github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU= github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= @@ -84,8 +78,8 @@ github.com/quic-go/qpack v0.6.0 h1:g7W+BMYynC1LbYLSqRt8PBg5Tgwxn214ZZR34VIOjz8= github.com/quic-go/qpack v0.6.0/go.mod h1:lUpLKChi8njB4ty2bFLX2x4gzDqXwUpaO1DP9qMDZII= github.com/quic-go/quic-go v0.60.0 h1:xcQioE8OM66UQLeUMHltK1CCcOu3JbVB4JAQdDQSB+0= github.com/quic-go/quic-go v0.60.0/go.mod h1:wpKpjmPpftl30sL6pFh7REVpjbcCVy4zt2vDyK1TuJk= -github.com/roadrunner-server/api-go/v6 v6.0.0-beta.12 h1:FcRcCvW9OfQvH45SFsI21VoHpOOov56OvOSnO4UKvXs= -github.com/roadrunner-server/api-go/v6 v6.0.0-beta.12/go.mod h1:prGWJ2GoF5YD5PIG7Tb6VKulU3bWoFwr9DCwgxheb80= +github.com/roadrunner-server/api-go/v6 v6.0.0-beta.12.0.20260714200341-93604e5012d4 h1:G5RlEP+rKKdarSw/ZcpWlpyrCne1AbuSZB8p9TagU3I= +github.com/roadrunner-server/api-go/v6 v6.0.0-beta.12.0.20260714200341-93604e5012d4/go.mod h1:Y4rsabWjr4Y10Jg6H8J5NDitQqlnXmGhCdgR+zyLYkI= github.com/roadrunner-server/api-plugins/v6 v6.0.0-beta.2 h1:GqsZzWQ5jMXRF1O/b8IqFz9PLpS7Ui0K4OyACLql2MI= github.com/roadrunner-server/api-plugins/v6 v6.0.0-beta.2/go.mod h1:2v4yUK5Kvbvq8C3IkDoBkuamq9h+7i/JLjyf7k1j5JM= github.com/roadrunner-server/config/v6 v6.0.0-beta.3 h1:G0EUzJ6Yw4UnleM6BhnOBbYPXKDHRmCJiGhC3nXDBwI= @@ -98,26 +92,26 @@ github.com/roadrunner-server/errors v1.5.0 h1:unG7LKIZrSzkCCF3YLRLA5VyqE0KKomofX github.com/roadrunner-server/errors v1.5.0/go.mod h1:g9fo/T2C13cWRDR9PW1r0ZAOSQfNhWAZawyfkGiaHuI= github.com/roadrunner-server/events v1.0.1 h1:waCkKhxhzdK3VcI1xG22l+h+0J+Nfdpxjhyy01Un+kI= github.com/roadrunner-server/events v1.0.1/go.mod h1:WZRqoEVaFm209t52EuoT7ISUtvX6BrCi6bI/7pjkVC0= -github.com/roadrunner-server/goridge/v4 v4.0.0-beta.2 h1:MgH6oiSgcl+vphsQ6JpyedkXQ/DPf8zVpn0z7rdBp10= -github.com/roadrunner-server/goridge/v4 v4.0.0-beta.2/go.mod h1:Wv9CBO9VIU92e5iZIuehLHKakXgMkOzxoT4/oHDjIUA= -github.com/roadrunner-server/http/v6 v6.0.0-beta.7 h1:uCKQBlD5gOCuGlJQA4h4q+IDK5C3tSGBPCPTUJjRztY= -github.com/roadrunner-server/http/v6 v6.0.0-beta.7/go.mod h1:T5XNZsqAsUMozEthOD5PmRhQYl7HYx7JqV7NAuA459Q= -github.com/roadrunner-server/informer/v6 v6.0.0-beta.2 h1:tJsNgbQ28mK5CdQCpU+BY6ScWP884nhpGYfwJalZOlU= -github.com/roadrunner-server/informer/v6 v6.0.0-beta.2/go.mod h1:nDn5jjR1ZI1Xhz01j32bSs4PHv78IE96rcqdFr/ZvgU= -github.com/roadrunner-server/jobs/v6 v6.0.0-beta.7 h1:RNb8fkVk2GO02E3FFyCmYz6WFkeXrBfV/daK84WXAD0= -github.com/roadrunner-server/jobs/v6 v6.0.0-beta.7/go.mod h1:kHTJXgjMe/7t9vJN4imEO3RlB7crnnp59X8TmgiEOqw= -github.com/roadrunner-server/kv/v6 v6.0.0-beta.5 h1:/Da5S0H72t1F/pzmgweT1rpZLREXOq8X7coIONEy/qY= -github.com/roadrunner-server/kv/v6 v6.0.0-beta.5/go.mod h1:OnALpec1My3YPZuV6/bawaGp8QdqY0M7M7rmoVeuYu0= +github.com/roadrunner-server/goridge/v4 v4.0.0-beta.2.0.20260714195909-75e9ece43063 h1:0mNGmXgYR2/hUhPQde+GJFzj/UQ6vqrLAHcwa5C5rqA= +github.com/roadrunner-server/goridge/v4 v4.0.0-beta.2.0.20260714195909-75e9ece43063/go.mod h1:1aHppV68y/VqRED/AsfNg59sft9aQOhqgr5Z5n49jbM= +github.com/roadrunner-server/http/v6 v6.0.0-beta.7.0.20260714202301-3d2c74eb5e61 h1:cfjEBPvn6M82SBqzQribvddESFxqb0AYyuMQRQ3mTjs= +github.com/roadrunner-server/http/v6 v6.0.0-beta.7.0.20260714202301-3d2c74eb5e61/go.mod h1:DpJJOUB9VuCYLM58LOzJuRAbiYZrbg1zNdEBdSzMQBg= +github.com/roadrunner-server/informer/v6 v6.0.0-beta.2.0.20260714201850-2854c943433a h1:AiM0l5YReltbsn+/LIS4MAfTD+XWSnijeIlffLh9q58= +github.com/roadrunner-server/informer/v6 v6.0.0-beta.2.0.20260714201850-2854c943433a/go.mod h1:uQGHMxNdwu85+i27kDMGTCVWsvBqQIZhTSGakFMSYu8= +github.com/roadrunner-server/jobs/v6 v6.0.0-beta.7.0.20260714202317-a904df360b84 h1:FW3BxRdj3Sz5aYQ/joDQWTb66KQqyqyu8mbrltX0s4k= +github.com/roadrunner-server/jobs/v6 v6.0.0-beta.7.0.20260714202317-a904df360b84/go.mod h1:MidnjQZiP9m3LfLJKtUlXZ75cU+xHuOpLap89jjnPIo= +github.com/roadrunner-server/kv/v6 v6.0.0-beta.5.0.20260714202229-7ef2a556a8ad h1:bFL6Cvs1jClAdfayB5whcPx++iTCBb2MDFpI/aISQ98= +github.com/roadrunner-server/kv/v6 v6.0.0-beta.5.0.20260714202229-7ef2a556a8ad/go.mod h1:4KXG/voKqBE3ev1iK/AYYqqrYxZLLcN7LDXfFvSj7BY= github.com/roadrunner-server/logger/v6 v6.0.0-beta.3 h1:eoJKXAUSyykDfVX6eTUhmAn6Y8pS/LyI5fDP4H+G5rQ= github.com/roadrunner-server/logger/v6 v6.0.0-beta.3/go.mod h1:MwHb3AbltHYtu7nRpml5NeYu7O+W8rCpDBeNTTEoE1M= github.com/roadrunner-server/pool/v2 v2.0.0-beta.1 h1:jpYXFtdD6QGAdAGPgMxrNi3j1CegCRpb2y+A+3GnXFA= github.com/roadrunner-server/pool/v2 v2.0.0-beta.1/go.mod h1:Bo1wT7RtL3eyQHXBUohNhtj/yAmRt6Rq8smuBg5pWkY= github.com/roadrunner-server/priority_queue v1.0.6 h1:x8bcMyjWs2Z4ySbO9BTP8Dzy2prCuazJY9HHrVTmUVY= github.com/roadrunner-server/priority_queue v1.0.6/go.mod h1:aJ2D9s18+OGpFfNgwoIduraaFYBGv4FKElnpzqO+TBI= -github.com/roadrunner-server/resetter/v6 v6.0.0-beta.3 h1:+hkbf/kXpvFjx4LfkuH8dvR07rBwrStmcqJaffsfL0g= -github.com/roadrunner-server/resetter/v6 v6.0.0-beta.3/go.mod h1:zV+MfVo6jtvrop+04HNcr4z3b/22qyKXukK29kzagYc= -github.com/roadrunner-server/rpc/v6 v6.0.0-beta.4 h1:Qj2nrHIWOHE9Tys+FBG2IdoPtzgIUh6juQ5wXLGGDMw= -github.com/roadrunner-server/rpc/v6 v6.0.0-beta.4/go.mod h1:k5KT3fpnJVd27m0HbGGBiTPXlWI6eJdd6C+ohp5IE0U= +github.com/roadrunner-server/resetter/v6 v6.0.0-beta.3.0.20260714201533-91a174bb65d4 h1:CTtuTFi2YYCkQx/Tu3ANfCPAkgle4AI6RGue70/Vb10= +github.com/roadrunner-server/resetter/v6 v6.0.0-beta.3.0.20260714201533-91a174bb65d4/go.mod h1:74CsBUTfnP9QhUpsUg2XqgurxCQYhxtPikk1C1iEfRc= +github.com/roadrunner-server/rpc/v6 v6.0.0-beta.4.0.20260714200548-15b82bc47898 h1:nc1MwAAG02mwOm76TO7sZZLt6ToKlkIz98shqAsxGoA= +github.com/roadrunner-server/rpc/v6 v6.0.0-beta.4.0.20260714200548-15b82bc47898/go.mod h1:pC636ll86dk1PJ2LoYbvurKMUARbF8Ckn4WBGvPXO8A= github.com/roadrunner-server/server/v6 v6.0.0-beta.6 h1:CPtH4eIYkeRKi5cPXxb0+J+LI824cqhIGXAfcH+nkjA= github.com/roadrunner-server/server/v6 v6.0.0-beta.6/go.mod h1:SbODuCzC2gcbFhAmJDWvjf34pPrUWP5NxxVsTRQDuZ4= github.com/roadrunner-server/tcplisten v1.5.2 h1:nn8yXYrhRDkfQ9AAu4V075uT4fZRmOnpxkawgE+bWPA= @@ -198,14 +192,8 @@ golang.org/x/text v0.38.0 h1:sXmwo9DwP3OK9EZ7PqAdaooSGozfl/3a6/xJcbzPRhE= golang.org/x/text v0.38.0/go.mod h1:YXZt3QhHUKYT53r2lLKFIVi6Ao1jdzrTR/KQ09qyxF4= golang.org/x/tools v0.47.0 h1:7Kn5x/d1svx/PzryTsqeoZN4TZwqeH5pGWjefhLi/1Q= golang.org/x/tools v0.47.0/go.mod h1:dFHnyTvFWY212G+h7ZY4Vsp/K3U4/7W9TyVaAul8uCA= -gonum.org/v1/gonum v0.17.0 h1:VbpOemQlsSMrYmn7T2OUvQ4dqxQXU+ouZFQsZOx50z4= -gonum.org/v1/gonum v0.17.0/go.mod h1:El3tOrEuMpv2UdMrbNlKEh9vd86bmQ6vqIcDwxEOc1E= google.golang.org/genproto v0.0.0-20260630182238-925bb5da69e7 h1:lQG76ePMKmtujel4VIVMiFoHVWVNtJdawbCZJtWlVXU= google.golang.org/genproto v0.0.0-20260630182238-925bb5da69e7/go.mod h1:LwlOWYBU335L+sR55UuR5fbbU8KmEX+3tUHf3SwMmhM= -google.golang.org/genproto/googleapis/rpc v0.0.0-20260630182238-925bb5da69e7 h1:eM/YSd5bBFagF51o1E745Ta7RwzpW0h+z+QDNZOgmQ8= -google.golang.org/genproto/googleapis/rpc v0.0.0-20260630182238-925bb5da69e7/go.mod h1:4Hqkh8ycfw05ld/3BWL7rJOSfebL2Q+DVDeRgYgxUU8= -google.golang.org/grpc v1.82.0 h1:vguDnZUPjE26w09A63VoxZPnvPjB5Riyc0mkXPFmAIU= -google.golang.org/grpc v1.82.0/go.mod h1:yzTZ1TB1Z3SG+LIYaI+WiE8D5+PZ3ArnrSp8zF3+/ZA= google.golang.org/protobuf v1.36.11 h1:fV6ZwhNocDyBLK0dj+fg8ektcVegBBuEolpbTQyBNVE= google.golang.org/protobuf v1.36.11/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= diff --git a/tests/helpers/helpers.go b/tests/helpers/helpers.go index 0ae9d4d..d6be1f7 100644 --- a/tests/helpers/helpers.go +++ b/tests/helpers/helpers.go @@ -1,67 +1,80 @@ package helpers import ( - "context" - "crypto/tls" "net" - "net/http" + "net/rpc" "slices" "testing" "time" - "connectrpc.com/connect" "github.com/google/uuid" jobsProto "github.com/roadrunner-server/api-go/v6/jobs/v2" - "github.com/roadrunner-server/api-go/v6/jobs/v2/jobsV2connect" - "github.com/roadrunner-server/api-go/v6/kv/v2/kvV2connect" jobState "github.com/roadrunner-server/api-plugins/v6/jobs" + goridgeRpc "github.com/roadrunner-server/goridge/v4/pkg/rpc" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" - "golang.org/x/net/http2" "google.golang.org/protobuf/types/known/emptypb" ) -func newHTTPClient(t *testing.T) *http.Client { +const ( + push = "jobs.Push" + pause = "jobs.Pause" + destroy = "jobs.Destroy" + resume = "jobs.Resume" + stat = "jobs.GetStats" +) + +// dialRPC opens a goridge net/rpc client against the RoadRunner RPC endpoint. +func dialRPC(t *testing.T, address string) *rpc.Client { t.Helper() - httpc := &http.Client{Transport: &http2.Transport{ - AllowHTTP: true, - DialTLSContext: func(ctx context.Context, network, addr string, _ *tls.Config) (net.Conn, error) { - return new(net.Dialer).DialContext(ctx, network, addr) - }, - }} - t.Cleanup(httpc.CloseIdleConnections) - return httpc + conn, err := (&net.Dialer{}).DialContext(t.Context(), "tcp", address) + require.NoError(t, err) + return rpc.NewClientWithCodec(goridgeRpc.NewClientCodec(conn)) } -func NewJobsClient(t *testing.T, address string) jobsV2connect.JobsServiceClient { +// NewJobsClient returns a goridge net/rpc client for the jobs plugin RPC surface. +// The connection is closed on test cleanup. +func NewJobsClient(t *testing.T, address string) *rpc.Client { t.Helper() - return jobsV2connect.NewJobsServiceClient(newHTTPClient(t), "http://"+address) + client := dialRPC(t, address) + t.Cleanup(func() { _ = client.Close() }) + return client } -func NewKVClient(t *testing.T, address string) kvV2connect.KvServiceClient { +// NewKVClient returns a goridge net/rpc client for the kv plugin RPC surface. +// The connection is closed on test cleanup. +func NewKVClient(t *testing.T, address string) *rpc.Client { t.Helper() - return kvV2connect.NewKvServiceClient(newHTTPClient(t), "http://"+address) + client := dialRPC(t, address) + t.Cleanup(func() { _ = client.Close() }) + return client } func ResumePipes(address string, pipes ...string) func(t *testing.T) { return func(t *testing.T) { - client := NewJobsClient(t, address) - _, err := client.Resume(t.Context(), connect.NewRequest(&jobsProto.Pipelines{Pipelines: slices.Clone(pipes)})) + client := dialRPC(t, address) + defer func() { _ = client.Close() }() + + err := client.Call(resume, &jobsProto.Pipelines{Pipelines: slices.Clone(pipes)}, &jobsProto.JobsHandlerResponse{}) require.NoError(t, err) } } func PushToPipe(pipeline string, autoAck bool, address string) func(t *testing.T) { return func(t *testing.T) { - client := NewJobsClient(t, address) - _, err := client.Push(t.Context(), connect.NewRequest(&jobsProto.PushRequest{Job: createDummyJob(pipeline, autoAck)})) + client := dialRPC(t, address) + defer func() { _ = client.Close() }() + + err := client.Call(push, &jobsProto.PushRequest{Job: createDummyJob(pipeline, autoAck)}, &jobsProto.JobsHandlerResponse{}) require.NoError(t, err) } } func PushToDisabledPipe(address, pipeline string) func(t *testing.T) { return func(t *testing.T) { - client := NewJobsClient(t, address) + client := dialRPC(t, address) + defer func() { _ = client.Close() }() + req := &jobsProto.PushRequest{Job: &jobsProto.Job{ Job: "some/php/namespace", Id: "1", @@ -72,14 +85,16 @@ func PushToDisabledPipe(address, pipeline string) func(t *testing.T) { Pipeline: pipeline, }, }} - _, err := client.Push(t.Context(), connect.NewRequest(req)) + err := client.Call(push, req, &jobsProto.JobsHandlerResponse{}) require.NoError(t, err) } } func PushToPipeDelayed(address string, pipeline string, delay int64) func(t *testing.T) { return func(t *testing.T) { - client := NewJobsClient(t, address) + client := dialRPC(t, address) + defer func() { _ = client.Close() }() + req := &jobsProto.PushRequest{Job: &jobsProto.Job{ Job: "some/php/namespace", Id: uuid.NewString(), @@ -91,7 +106,7 @@ func PushToPipeDelayed(address string, pipeline string, delay int64) func(t *tes Delay: delay, }, }} - _, err := client.Push(t.Context(), connect.NewRequest(req)) + err := client.Call(push, req, &jobsProto.JobsHandlerResponse{}) assert.NoError(t, err) } } @@ -113,22 +128,26 @@ func createDummyJob(pipeline string, autoAck bool) *jobsProto.Job { func PausePipelines(address string, pipes ...string) func(t *testing.T) { return func(t *testing.T) { - client := NewJobsClient(t, address) - _, err := client.Pause(t.Context(), connect.NewRequest(&jobsProto.Pipelines{Pipelines: slices.Clone(pipes)})) + client := dialRPC(t, address) + defer func() { _ = client.Close() }() + + err := client.Call(pause, &jobsProto.Pipelines{Pipelines: slices.Clone(pipes)}, &jobsProto.JobsHandlerResponse{}) assert.NoError(t, err) } } func DestroyPipelines(address string, pipes ...string) func(t *testing.T) { return func(t *testing.T) { - client := NewJobsClient(t, address) + client := dialRPC(t, address) + defer func() { _ = client.Close() }() + req := &jobsProto.Pipelines{Pipelines: slices.Clone(pipes)} // Retry the destroy 10× with 1s gaps; if all attempts fail, return // without asserting. Some negative tests intentionally destroy // non-existent pipelines and rely on this silent-after-retry pattern. for range 10 { - _, err := client.Destroy(t.Context(), connect.NewRequest(req)) + err := client.Call(destroy, req, &jobsProto.Pipelines{}) if err == nil { return } @@ -139,21 +158,22 @@ func DestroyPipelines(address string, pipes ...string) func(t *testing.T) { func Stats(address string, state *jobState.State) func(t *testing.T) { return func(t *testing.T) { - client := NewJobsClient(t, address) + client := dialRPC(t, address) + defer func() { _ = client.Close() }() - resp, err := client.GetStats(t.Context(), connect.NewRequest(&emptypb.Empty{})) + st := &jobsProto.Stats{} + err := client.Call(stat, &emptypb.Empty{}, st) require.NoError(t, err) - require.NotNil(t, resp) - require.NotEmpty(t, resp.Msg.GetStats()) - - st := resp.Msg.GetStats()[0] - state.Queue = st.GetQueue() - state.Pipeline = st.GetPipeline() - state.Driver = st.GetDriver() - state.Active = st.GetActive() - state.Delayed = st.GetDelayed() - state.Reserved = st.GetReserved() - state.Ready = st.GetReady() - state.Priority = st.GetPriority() + require.NotEmpty(t, st.GetStats()) + + s := st.GetStats()[0] + state.Queue = s.GetQueue() + state.Pipeline = s.GetPipeline() + state.Driver = s.GetDriver() + state.Active = s.GetActive() + state.Delayed = s.GetDelayed() + state.Reserved = s.GetReserved() + state.Ready = s.GetReady() + state.Priority = s.GetPriority() } } diff --git a/tests/jobs_memory_test.go b/tests/jobs_memory_test.go index 0eed70e..4355ead 100644 --- a/tests/jobs_memory_test.go +++ b/tests/jobs_memory_test.go @@ -14,7 +14,6 @@ import ( "tests/helpers" mocklogger "tests/mock" - "connectrpc.com/connect" 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" @@ -1008,7 +1007,7 @@ func declareMemoryPipe(prefetch string) func(t *testing.T) { "prefetch": prefetch, "priority": "33", }} - _, err := client.Declare(t.Context(), connect.NewRequest(req)) + err := client.Call("jobs.Declare", req, &jobsProto.JobsHandlerResponse{}) assert.NoError(t, err) } } @@ -1016,7 +1015,7 @@ func declareMemoryPipe(prefetch string) func(t *testing.T) { func consumeMemoryPipe(pipelines []string) func(t *testing.T) { return func(t *testing.T) { client := helpers.NewJobsClient(t, "127.0.0.1:6001") - _, err := client.Resume(t.Context(), connect.NewRequest(&jobsProto.Pipelines{Pipelines: slices.Clone(pipelines)})) + err := client.Call("jobs.Resume", &jobsProto.Pipelines{Pipelines: slices.Clone(pipelines)}, &jobsProto.JobsHandlerResponse{}) assert.NoError(t, err) } } diff --git a/tests/jobs_native_test.go b/tests/jobs_native_test.go deleted file mode 100644 index bab5835..0000000 --- a/tests/jobs_native_test.go +++ /dev/null @@ -1,51 +0,0 @@ -package memory - -import ( - "context" - "net/http" - "net/http/httptest" - "testing" - - "connectrpc.com/connect" - jobsProto "github.com/roadrunner-server/api-go/v6/jobs/v2" - "github.com/roadrunner-server/api-go/v6/jobs/v2/jobsV2connect" - "github.com/stretchr/testify/require" -) - -// fakeJobsService stubs only Declare — the procedure the PHP worker -// `jobs_create_memory.php` would have invoked via spiral/goridge to register -// a new in-memory pipeline at runtime. Other methods fall through to -// UnimplementedJobsServiceHandler (CodeUnimplemented). -type fakeJobsService struct { - jobsV2connect.UnimplementedJobsServiceHandler -} - -func (fakeJobsService) Declare( - _ context.Context, _ *connect.Request[jobsProto.DeclareRequest], -) (*connect.Response[jobsProto.JobsHandlerResponse], error) { - return connect.NewResponse(&jobsProto.JobsHandlerResponse{}), nil -} - -// TestJobsNativeDeclare is a pure request/response Connect-RPC smoke test for -// jobs.JobsService.Declare, mirroring what the (still-broken) PHP -// TestMemoryCreate exercises via Jobs.create(MemoryCreateInfo). No -// Roadrunner container, no PHP — just proves the proto types + connectrpc -// wire round-trip for this procedure. -func TestJobsNativeDeclare(t *testing.T) { - mux := http.NewServeMux() - mux.Handle(jobsV2connect.NewJobsServiceHandler(fakeJobsService{})) - - srv := httptest.NewServer(mux) - t.Cleanup(srv.Close) - - client := jobsV2connect.NewJobsServiceClient(srv.Client(), srv.URL) - _, err := client.Declare(t.Context(), connect.NewRequest(&jobsProto.DeclareRequest{ - Pipeline: map[string]string{ - "driver": "memory", - "name": "example", - "priority": "10", - "prefetch": "100", - }, - })) - require.NoError(t, err) -} diff --git a/tests/kv_memory_test.go b/tests/kv_memory_test.go index 02ad6d9..2cc0675 100644 --- a/tests/kv_memory_test.go +++ b/tests/kv_memory_test.go @@ -14,7 +14,6 @@ import ( "tests/helpers" - "connectrpc.com/connect" kvProto "github.com/roadrunner-server/api-go/v6/kv/v2" "github.com/roadrunner-server/config/v6" "github.com/roadrunner-server/endure/v2" @@ -177,7 +176,6 @@ func TestSetManyMemory(t *testing.T) { ngprev := runtime.NumGoroutine() client := helpers.NewKVClient(t, "127.0.0.1:6666") - ctx := t.Context() tt := durationpb.New(time.Minute * 10) data := &kvProto.KvRequest{ @@ -191,7 +189,7 @@ func TestSetManyMemory(t *testing.T) { } for range 10_000 { - _, err := client.Set(ctx, connect.NewRequest(data)) + err := client.Call("kv.Set", data, &kvProto.KvResponse{}) require.NoError(t, err) } runtime.GC() @@ -215,7 +213,7 @@ func TestSetManyMemory(t *testing.T) { time.Sleep(time.Second * 5) - _, err = client.Clear(ctx, connect.NewRequest(data)) + err = client.Call("kv.Clear", data, &kvProto.KvResponse{}) require.NoError(t, err) stopCh <- struct{}{} @@ -289,7 +287,6 @@ func testRPCMethodsInMemory(t *testing.T) { const storage = "memory-rr" client := helpers.NewKVClient(t, "127.0.0.1:6001") - ctx := t.Context() tt := durationpb.New(time.Second * 5) keys := &kvProto.KvRequest{ @@ -312,23 +309,26 @@ func testRPCMethodsInMemory(t *testing.T) { }, } - _, err := client.Set(ctx, connect.NewRequest(data)) + err := client.Call("kv.Set", data, &kvProto.KvResponse{}) assert.NoError(t, err) - resp, err := client.Has(ctx, connect.NewRequest(keys)) + resp := &kvProto.KvResponse{} + err = client.Call("kv.Has", keys, resp) assert.NoError(t, err) - assert.Len(t, resp.Msg.GetItems(), 3) + assert.Len(t, resp.GetItems(), 3) // key "c" should be deleted time.Sleep(time.Second * 7) - resp, err = client.Has(ctx, connect.NewRequest(keys)) + resp = &kvProto.KvResponse{} + err = client.Call("kv.Has", keys, resp) assert.NoError(t, err) - assert.Len(t, resp.Msg.GetItems(), 2) + assert.Len(t, resp.GetItems(), 2) - resp, err = client.MGet(ctx, connect.NewRequest(keys)) + resp = &kvProto.KvResponse{} + err = client.Call("kv.MGet", keys, resp) assert.NoError(t, err) - assert.Len(t, resp.Msg.GetItems(), 2) // c is expired + assert.Len(t, resp.GetItems(), 2) // c is expired tt2 := durationpb.New(time.Second * 10) @@ -341,7 +341,7 @@ func testRPCMethodsInMemory(t *testing.T) { }, } - _, err = client.MExpire(ctx, connect.NewRequest(data2)) + err = client.Call("kv.MExpire", data2, &kvProto.KvResponse{}) assert.NoError(t, err) keys2 := &kvProto.KvRequest{ @@ -353,27 +353,30 @@ func testRPCMethodsInMemory(t *testing.T) { }, } - resp, err = client.TTL(ctx, connect.NewRequest(keys2)) + resp = &kvProto.KvResponse{} + err = client.Call("kv.TTL", keys2, resp) assert.NoError(t, err) - assert.Len(t, resp.Msg.GetItems(), 3) + assert.Len(t, resp.GetItems(), 3) // HAS AFTER TTL time.Sleep(time.Second * 15) - resp, err = client.Has(ctx, connect.NewRequest(keys2)) + resp = &kvProto.KvResponse{} + err = client.Call("kv.Has", keys2, resp) assert.NoError(t, err) - assert.Empty(t, resp.Msg.GetItems()) + assert.Empty(t, resp.GetItems()) keysDel := &kvProto.KvRequest{ Storage: storage, Items: []*kvProto.KvItem{{Key: "e"}}, } - _, err = client.Delete(ctx, connect.NewRequest(keysDel)) + err = client.Call("kv.Delete", keysDel, &kvProto.KvResponse{}) assert.NoError(t, err) - resp, err = client.Has(ctx, connect.NewRequest(keysDel)) + resp = &kvProto.KvResponse{} + err = client.Call("kv.Has", keysDel, resp) assert.NoError(t, err) - assert.Empty(t, resp.Msg.GetItems()) + assert.Empty(t, resp.GetItems()) dataClear := &kvProto.KvRequest{ Storage: storage, @@ -386,21 +389,23 @@ func testRPCMethodsInMemory(t *testing.T) { }, } - _, err = client.Set(ctx, connect.NewRequest(dataClear)) + err = client.Call("kv.Set", dataClear, &kvProto.KvResponse{}) assert.NoError(t, err) - resp, err = client.Has(ctx, connect.NewRequest(dataClear)) + resp = &kvProto.KvResponse{} + err = client.Call("kv.Has", dataClear, resp) assert.NoError(t, err) - assert.Len(t, resp.Msg.GetItems(), 5) + assert.Len(t, resp.GetItems(), 5) - _, err = client.Clear(ctx, connect.NewRequest(&kvProto.KvRequest{Storage: storage})) + err = client.Call("kv.Clear", &kvProto.KvRequest{Storage: storage}, &kvProto.KvResponse{}) assert.NoError(t, err) - resp, err = client.Has(ctx, connect.NewRequest(dataClear)) + resp = &kvProto.KvResponse{} + err = client.Call("kv.Has", dataClear, resp) assert.NoError(t, err) - assert.Empty(t, resp.Msg.GetItems()) + assert.Empty(t, resp.GetItems()) - _, err = client.Clear(ctx, connect.NewRequest(data)) + err = client.Call("kv.Clear", data, &kvProto.KvResponse{}) require.NoError(t, err) } @@ -467,7 +472,6 @@ func TestInMemoryKVTracer(t *testing.T) { const storage = "memory-rr" client := helpers.NewKVClient(t, "127.0.0.1:6001") - ctx := t.Context() tt := durationpb.New(time.Second * 30) @@ -478,42 +482,45 @@ func TestInMemoryKVTracer(t *testing.T) { {Key: "b", Value: []byte("bb")}, }, } - _, err = client.Set(ctx, connect.NewRequest(data)) + err = client.Call("kv.Set", data, &kvProto.KvResponse{}) assert.NoError(t, err) keys := &kvProto.KvRequest{ Storage: storage, Items: []*kvProto.KvItem{{Key: "a"}, {Key: "b"}}, } - resp, err := client.Has(ctx, connect.NewRequest(keys)) + resp := &kvProto.KvResponse{} + err = client.Call("kv.Has", keys, resp) assert.NoError(t, err) - assert.Len(t, resp.Msg.GetItems(), 2) + assert.Len(t, resp.GetItems(), 2) - resp, err = client.MGet(ctx, connect.NewRequest(keys)) + resp = &kvProto.KvResponse{} + err = client.Call("kv.MGet", keys, resp) assert.NoError(t, err) - assert.Len(t, resp.Msg.GetItems(), 2) + assert.Len(t, resp.GetItems(), 2) - resp, err = client.TTL(ctx, connect.NewRequest(&kvProto.KvRequest{ + resp = &kvProto.KvResponse{} + err = client.Call("kv.TTL", &kvProto.KvRequest{ Storage: storage, Items: []*kvProto.KvItem{{Key: "a"}}, - })) + }, resp) assert.NoError(t, err) - assert.Len(t, resp.Msg.GetItems(), 1) + assert.Len(t, resp.GetItems(), 1) tt2 := durationpb.New(time.Second * 60) - _, err = client.MExpire(ctx, connect.NewRequest(&kvProto.KvRequest{ + err = client.Call("kv.MExpire", &kvProto.KvRequest{ Storage: storage, Items: []*kvProto.KvItem{{Key: "b", Ttl: tt2}}, - })) + }, &kvProto.KvResponse{}) assert.NoError(t, err) - _, err = client.Delete(ctx, connect.NewRequest(&kvProto.KvRequest{ + err = client.Call("kv.Delete", &kvProto.KvRequest{ Storage: storage, Items: []*kvProto.KvItem{{Key: "b"}}, - })) + }, &kvProto.KvResponse{}) assert.NoError(t, err) - _, err = client.Clear(ctx, connect.NewRequest(&kvProto.KvRequest{Storage: storage})) + err = client.Call("kv.Clear", &kvProto.KvRequest{Storage: storage}, &kvProto.KvResponse{}) assert.NoError(t, err) stopCh <- struct{}{} diff --git a/tests/kv_native_test.go b/tests/kv_native_test.go deleted file mode 100644 index c6cde9b..0000000 --- a/tests/kv_native_test.go +++ /dev/null @@ -1,49 +0,0 @@ -package memory - -import ( - "context" - "net/http" - "net/http/httptest" - "testing" - - "connectrpc.com/connect" - kvProto "github.com/roadrunner-server/api-go/v6/kv/v2" - "github.com/roadrunner-server/api-go/v6/kv/v2/kvV2connect" - "github.com/stretchr/testify/require" -) - -// fakeKvService stubs only Has — the procedure the PHP worker -// `kv-order.php` would have invoked via spiral/goridge to probe whether -// "test_key" is present in the "memory-rr" storage. Other methods fall -// through to UnimplementedKvServiceHandler (CodeUnimplemented). -type fakeKvService struct { - kvV2connect.UnimplementedKvServiceHandler -} - -func (fakeKvService) Has( - _ context.Context, _ *connect.Request[kvProto.KvRequest], -) (*connect.Response[kvProto.KvResponse], error) { - // Empty Items mirrors what an empty memory storage would have returned - // for the PHP worker's startup get("test_key") probe. - return connect.NewResponse(&kvProto.KvResponse{}), nil -} - -// TestKVNativeHas is a pure request/response Connect-RPC smoke test for -// kv.KvService.Has, mirroring what the (still-broken) PHP TestInMemoryOrder -// path exercises. No Roadrunner container, no PHP — just proves the proto -// types + connectrpc wire round-trip for this procedure. -func TestKVNativeHas(t *testing.T) { - mux := http.NewServeMux() - mux.Handle(kvV2connect.NewKvServiceHandler(fakeKvService{})) - - srv := httptest.NewServer(mux) - t.Cleanup(srv.Close) - - client := kvV2connect.NewKvServiceClient(srv.Client(), srv.URL) - resp, err := client.Has(t.Context(), connect.NewRequest(&kvProto.KvRequest{ - Storage: "memory-rr", - Items: []*kvProto.KvItem{{Key: "test_key"}}, - })) - require.NoError(t, err) - require.Empty(t, resp.Msg.GetItems()) -}