diff --git a/go.mod b/go.mod index baa0ecb..21554ae 100644 --- a/go.mod +++ b/go.mod @@ -4,16 +4,6 @@ go 1.26 toolchain go1.26.4 -require ( - connectrpc.com/connect v1.20.0 - github.com/roadrunner-server/api-go/v6 v6.0.0-beta.12 -) +require github.com/roadrunner-server/api-go/v6 v6.0.0-beta.12.0.20260714200341-93604e5012d4 -require ( - golang.org/x/net v0.56.0 // indirect - golang.org/x/sys v0.46.0 // indirect - golang.org/x/text v0.38.0 // indirect - google.golang.org/genproto/googleapis/rpc v0.0.0-20260630182238-925bb5da69e7 // indirect - google.golang.org/grpc v1.82.0 // indirect - google.golang.org/protobuf v1.36.11 // indirect -) +require google.golang.org/protobuf v1.36.11 // indirect diff --git a/go.sum b/go.sum index 8d0d064..213d39e 100644 --- a/go.sum +++ b/go.sum @@ -1,42 +1,6 @@ -connectrpc.com/connect v1.20.0 h1:6TNDAB+WeNd2uolWNlYczB5E0KNNaVMNUEx8JEUsPmQ= -connectrpc.com/connect v1.20.0/go.mod h1:A2ygJrukXwWy32vkCAAHNVguZrqZ+jeZ9rGRnGR4dN4= -github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= -github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= -github.com/go-logr/logr v1.4.3 h1:CjnDlHq8ikf6E492q6eKboGOC0T8CDaOvkHCIg8idEI= -github.com/go-logr/logr v1.4.3/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY= -github.com/go-logr/stdr v1.2.2 h1:hSWxHoqTgW2S2qGc0LTAI563KZ5YKYRhT3MFKZMbjag= -github.com/go-logr/stdr v1.2.2/go.mod h1:mMo/vtBO5dYbehREoey6XUKy/eSumjCCveDpRre4VKE= -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= -github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= -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= -go.opentelemetry.io/auto/sdk v1.2.1 h1:jXsnJ4Lmnqd11kwkBV2LgLoFMZKizbCi5fNZ/ipaZ64= -go.opentelemetry.io/auto/sdk v1.2.1/go.mod h1:KRTj+aOaElaLi+wW1kO/DZRXwkF4C5xPbEe3ZiIhN7Y= -go.opentelemetry.io/otel v1.43.0 h1:mYIM03dnh5zfN7HautFE4ieIig9amkNANT+xcVxAj9I= -go.opentelemetry.io/otel v1.43.0/go.mod h1:JuG+u74mvjvcm8vj8pI5XiHy1zDeoCS2LB1spIq7Ay0= -go.opentelemetry.io/otel/metric v1.43.0 h1:d7638QeInOnuwOONPp4JAOGfbCEpYb+K6DVWvdxGzgM= -go.opentelemetry.io/otel/metric v1.43.0/go.mod h1:RDnPtIxvqlgO8GRW18W6Z/4P462ldprJtfxHxyKd2PY= -go.opentelemetry.io/otel/sdk v1.43.0 h1:pi5mE86i5rTeLXqoF/hhiBtUNcrAGHLKQdhg4h4V9Dg= -go.opentelemetry.io/otel/sdk v1.43.0/go.mod h1:P+IkVU3iWukmiit/Yf9AWvpyRDlUeBaRg6Y+C58QHzg= -go.opentelemetry.io/otel/sdk/metric v1.43.0 h1:S88dyqXjJkuBNLeMcVPRFXpRw2fuwdvfCGLEo89fDkw= -go.opentelemetry.io/otel/sdk/metric v1.43.0/go.mod h1:C/RJtwSEJ5hzTiUz5pXF1kILHStzb9zFlIEe85bhj6A= -go.opentelemetry.io/otel/trace v1.43.0 h1:BkNrHpup+4k4w+ZZ86CZoHHEkohws8AY+WTX09nk+3A= -go.opentelemetry.io/otel/trace v1.43.0/go.mod h1:/QJhyVBUUswCphDVxq+8mld+AvhXZLhe+8WVFxiFff0= -golang.org/x/net v0.56.0 h1:Rw8j/hFzGvJUZwNBXnAtf5sVDVt+65SK2C7IxCxZt5o= -golang.org/x/net v0.56.0/go.mod h1:D3Ku6r+V6JROoZK144D2XfMHFcMq/0zSfLelVTCFKec= -golang.org/x/sys v0.46.0 h1:noSf2Fq6F8DBgS+LysIkx7rIExoNHJsxOAtPp4rthXw= -golang.org/x/sys v0.46.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= -golang.org/x/text v0.38.0 h1:sXmwo9DwP3OK9EZ7PqAdaooSGozfl/3a6/xJcbzPRhE= -golang.org/x/text v0.38.0/go.mod h1:YXZt3QhHUKYT53r2lLKFIVi6Ao1jdzrTR/KQ09qyxF4= -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/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= +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= google.golang.org/protobuf v1.36.11 h1:fV6ZwhNocDyBLK0dj+fg8ektcVegBBuEolpbTQyBNVE= google.golang.org/protobuf v1.36.11/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco= diff --git a/plugin.go b/plugin.go index b9b5ed6..d94b762 100644 --- a/plugin.go +++ b/plugin.go @@ -3,9 +3,6 @@ package lock import ( "context" "log/slog" - "net/http" - - "github.com/roadrunner-server/api-go/v6/lock/v1/lockV1connect" ) const pluginName string = "lock" @@ -43,6 +40,6 @@ func (p *Plugin) Name() string { return pluginName } -func (p *Plugin) RPC() (string, http.Handler) { - return lockV1connect.NewLockServiceHandler(&rpc{pl: p}) +func (p *Plugin) RPC() any { + return &rpc{pl: p} } diff --git a/rpc.go b/rpc.go index eede395..51ef7ce 100644 --- a/rpc.go +++ b/rpc.go @@ -5,7 +5,6 @@ import ( "errors" "time" - "connectrpc.com/connect" lockV1 "github.com/roadrunner-server/api-go/v6/lock/v1" ) @@ -24,88 +23,82 @@ func waitContext(parent context.Context, waitUs int64) (context.Context, context return context.WithTimeout(parent, time.Microsecond*time.Duration(waitUs)) } -func (r *rpc) Lock(ctx context.Context, req *connect.Request[lockV1.LockRequest]) (*connect.Response[lockV1.LockResponse], error) { - r.pl.log.Debug("lock request received", "ttl", int(req.Msg.GetTtl()), "wait_ttl", int(req.Msg.GetWait()), "resource", req.Msg.GetResource(), "id", req.Msg.GetId()) +func (r *rpc) Lock(in *lockV1.LockRequest, out *lockV1.LockResponse) error { + r.pl.log.Debug("lock request received", "ttl", int(in.GetTtl()), "wait_ttl", int(in.GetWait()), "resource", in.GetResource(), "id", in.GetId()) - if req.Msg.GetId() == "" { - return nil, connect.NewError(connect.CodeInvalidArgument, errEmptyID) + if in.GetId() == "" { + return errEmptyID } - cctx, cancel := waitContext(ctx, req.Msg.GetWait()) + cctx, cancel := waitContext(context.Background(), in.GetWait()) defer cancel() - return connect.NewResponse(&lockV1.LockResponse{ - Ok: r.pl.locks.lock(cctx, req.Msg.GetResource(), req.Msg.GetId(), int(req.Msg.GetTtl())), - }), nil + out.Ok = r.pl.locks.lock(cctx, in.GetResource(), in.GetId(), int(in.GetTtl())) + return nil } -func (r *rpc) LockRead(ctx context.Context, req *connect.Request[lockV1.LockRequest]) (*connect.Response[lockV1.LockResponse], error) { - r.pl.log.Debug("read lock request received", "ttl", int(req.Msg.GetTtl()), "wait_ttl", int(req.Msg.GetWait()), "resource", req.Msg.GetResource(), "id", req.Msg.GetId()) +func (r *rpc) LockRead(in *lockV1.LockRequest, out *lockV1.LockResponse) error { + r.pl.log.Debug("read lock request received", "ttl", int(in.GetTtl()), "wait_ttl", int(in.GetWait()), "resource", in.GetResource(), "id", in.GetId()) - if req.Msg.GetId() == "" { - return nil, connect.NewError(connect.CodeInvalidArgument, errEmptyID) + if in.GetId() == "" { + return errEmptyID } - cctx, cancel := waitContext(ctx, req.Msg.GetWait()) + cctx, cancel := waitContext(context.Background(), in.GetWait()) defer cancel() - return connect.NewResponse(&lockV1.LockResponse{ - Ok: r.pl.locks.lockRead(cctx, req.Msg.GetResource(), req.Msg.GetId(), int(req.Msg.GetTtl())), - }), nil + out.Ok = r.pl.locks.lockRead(cctx, in.GetResource(), in.GetId(), int(in.GetTtl())) + return nil } -func (r *rpc) Release(ctx context.Context, req *connect.Request[lockV1.LockRequest]) (*connect.Response[lockV1.LockResponse], error) { - r.pl.log.Debug("release request received", "ttl", int(req.Msg.GetTtl()), "wait_ttl", int(req.Msg.GetWait()), "resource", req.Msg.GetResource(), "id", req.Msg.GetId()) +func (r *rpc) Release(in *lockV1.LockRequest, out *lockV1.LockResponse) error { + r.pl.log.Debug("release request received", "ttl", int(in.GetTtl()), "wait_ttl", int(in.GetWait()), "resource", in.GetResource(), "id", in.GetId()) - if req.Msg.GetId() == "" { - return nil, connect.NewError(connect.CodeInvalidArgument, errEmptyID) + if in.GetId() == "" { + return errEmptyID } - cctx, cancel := waitContext(ctx, req.Msg.GetWait()) + cctx, cancel := waitContext(context.Background(), in.GetWait()) defer cancel() - return connect.NewResponse(&lockV1.LockResponse{ - Ok: r.pl.locks.release(cctx, req.Msg.GetResource(), req.Msg.GetId()), - }), nil + out.Ok = r.pl.locks.release(cctx, in.GetResource(), in.GetId()) + return nil } -func (r *rpc) ForceRelease(ctx context.Context, req *connect.Request[lockV1.LockRequest]) (*connect.Response[lockV1.LockResponse], error) { - r.pl.log.Debug("force release request received", "ttl", int(req.Msg.GetTtl()), "wait_ttl", int(req.Msg.GetWait()), "resource", req.Msg.GetResource(), "id", req.Msg.GetId()) +func (r *rpc) ForceRelease(in *lockV1.LockRequest, out *lockV1.LockResponse) error { + r.pl.log.Debug("force release request received", "ttl", int(in.GetTtl()), "wait_ttl", int(in.GetWait()), "resource", in.GetResource(), "id", in.GetId()) - cctx, cancel := waitContext(ctx, req.Msg.GetWait()) + cctx, cancel := waitContext(context.Background(), in.GetWait()) defer cancel() - return connect.NewResponse(&lockV1.LockResponse{ - Ok: r.pl.locks.forceRelease(cctx, req.Msg.GetResource()), - }), nil + out.Ok = r.pl.locks.forceRelease(cctx, in.GetResource()) + return nil } -func (r *rpc) Exists(ctx context.Context, req *connect.Request[lockV1.LockRequest]) (*connect.Response[lockV1.LockResponse], error) { - r.pl.log.Debug("exists request received", "ttl", int(req.Msg.GetTtl()), "wait_ttl", int(req.Msg.GetWait()), "resource", req.Msg.GetResource(), "id", req.Msg.GetId()) +func (r *rpc) Exists(in *lockV1.LockRequest, out *lockV1.LockResponse) error { + r.pl.log.Debug("exists request received", "ttl", int(in.GetTtl()), "wait_ttl", int(in.GetWait()), "resource", in.GetResource(), "id", in.GetId()) - if req.Msg.GetId() == "" { - return nil, connect.NewError(connect.CodeInvalidArgument, errEmptyID) + if in.GetId() == "" { + return errEmptyID } - cctx, cancel := waitContext(ctx, req.Msg.GetWait()) + cctx, cancel := waitContext(context.Background(), in.GetWait()) defer cancel() - return connect.NewResponse(&lockV1.LockResponse{ - Ok: r.pl.locks.exists(cctx, req.Msg.GetResource(), req.Msg.GetId()), - }), nil + out.Ok = r.pl.locks.exists(cctx, in.GetResource(), in.GetId()) + return nil } -func (r *rpc) UpdateTTL(ctx context.Context, req *connect.Request[lockV1.LockRequest]) (*connect.Response[lockV1.LockResponse], error) { - r.pl.log.Debug("updateTTL request received", "ttl", int(req.Msg.GetTtl()), "wait_ttl", int(req.Msg.GetWait()), "resource", req.Msg.GetResource(), "id", req.Msg.GetId()) +func (r *rpc) UpdateTTL(in *lockV1.LockRequest, out *lockV1.LockResponse) error { + r.pl.log.Debug("updateTTL request received", "ttl", int(in.GetTtl()), "wait_ttl", int(in.GetWait()), "resource", in.GetResource(), "id", in.GetId()) - if req.Msg.GetId() == "" { - return nil, connect.NewError(connect.CodeInvalidArgument, errEmptyID) + if in.GetId() == "" { + return errEmptyID } - cctx, cancel := waitContext(ctx, req.Msg.GetWait()) + cctx, cancel := waitContext(context.Background(), in.GetWait()) defer cancel() - return connect.NewResponse(&lockV1.LockResponse{ - Ok: r.pl.locks.updateTTL(cctx, req.Msg.GetResource(), req.Msg.GetId(), int(req.Msg.GetTtl())), - }), nil + out.Ok = r.pl.locks.updateTTL(cctx, in.GetResource(), in.GetId(), int(in.GetTtl())) + return nil } diff --git a/tests/configs/.rr-lock-api.yaml b/tests/configs/.rr-lock-api.yaml deleted file mode 100644 index 865cd05..0000000 --- a/tests/configs/.rr-lock-api.yaml +++ /dev/null @@ -1,9 +0,0 @@ -version: '3' - -rpc: - listen: tcp://127.0.0.1:6001 - -logs: - level: error - encoding: console - mode: development diff --git a/tests/go.mod b/tests/go.mod index 2379d68..625cf85 100644 --- a/tests/go.mod +++ b/tests/go.mod @@ -5,23 +5,19 @@ go 1.26 toolchain go1.26.4 require ( - connectrpc.com/connect v1.20.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/config/v6 v6.0.0-beta.3 github.com/roadrunner-server/endure/v2 v2.6.2 + github.com/roadrunner-server/goridge/v4 v4.0.0-beta.2.0.20260714195909-75e9ece43063 github.com/roadrunner-server/lock/v6 v6.0.0 github.com/roadrunner-server/logger/v6 v6.0.0-beta.3 - github.com/roadrunner-server/rpc/v6 v6.0.0-beta.4 + github.com/roadrunner-server/rpc/v6 v6.0.0-beta.4.0.20260714200548-15b82bc47898 github.com/stretchr/testify v1.11.1 - golang.org/x/net v0.56.0 - google.golang.org/grpc v1.82.0 - google.golang.org/protobuf v1.36.11 ) replace github.com/roadrunner-server/lock/v6 => ../ require ( - connectrpc.com/grpcreflect v1.3.0 // indirect github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc // indirect github.com/fatih/color v1.19.0 // indirect github.com/fsnotify/fsnotify v1.10.1 // indirect @@ -44,6 +40,6 @@ require ( go.yaml.in/yaml/v3 v3.0.4 // indirect golang.org/x/sys v0.46.0 // indirect golang.org/x/text v0.38.0 // indirect - google.golang.org/genproto/googleapis/rpc v0.0.0-20260630182238-925bb5da69e7 // indirect + google.golang.org/protobuf v1.36.11 // indirect gopkg.in/yaml.v3 v3.0.1 // indirect ) diff --git a/tests/go.sum b/tests/go.sum index 7884c53..1eeebe7 100644 --- a/tests/go.sum +++ b/tests/go.sum @@ -1,9 +1,3 @@ -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/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= -github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc h1:U9qPSI2PIWSS1VwoXQT9A3Wy9MM3WgvqSxFWenqJduM= github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/fatih/color v1.19.0 h1:Zp3PiM21/9Ld6FzSKyL5c/BULoe/ONr9KlbYVOfG8+w= @@ -12,18 +6,10 @@ github.com/frankban/quicktest v1.14.6 h1:7Xjx+VpznH+oBnejlPUj8oUpdxnVs4f8XU8WnHk github.com/frankban/quicktest v1.14.6/go.mod h1:4ptaffx2x8+WTWXmUCuVU6aPUX1/Mz7zb5vbUoiM6w0= github.com/fsnotify/fsnotify v1.10.1 h1:b0/UzAf9yR5rhf3RPm9gf3ehBPpf0oZKIjtpKrx59Ho= github.com/fsnotify/fsnotify v1.10.1/go.mod h1:TLheqan6HD6GBK6PrDWyDPBaEV8LspOxvPSjC+bVfgo= -github.com/go-logr/logr v1.4.3 h1:CjnDlHq8ikf6E492q6eKboGOC0T8CDaOvkHCIg8idEI= -github.com/go-logr/logr v1.4.3/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY= -github.com/go-logr/stdr v1.2.2 h1:hSWxHoqTgW2S2qGc0LTAI563KZ5YKYRhT3MFKZMbjag= -github.com/go-logr/stdr v1.2.2/go.mod h1:mMo/vtBO5dYbehREoey6XUKy/eSumjCCveDpRre4VKE= 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= -github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= github.com/joho/godotenv v1.5.1 h1:7eLL/+HRGLY0ldzfGMeQkb7vMd0as4CfYvUVzLqw0N0= github.com/joho/godotenv v1.5.1/go.mod h1:f4LDr5Voq0i2e/R5DDNOoa2zzDfwtkZa6DnEwAbqwq4= github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= @@ -38,18 +24,20 @@ github.com/pelletier/go-toml/v2 v2.4.2 h1:M2fKKbmyvI+hGId/D0W64qDBMVhJnNR10O5gIb github.com/pelletier/go-toml/v2 v2.4.2/go.mod h1:2gIqNv+qfxSVS7cM2xJQKtLSTLUE9V8t9Stt+h56mCY= github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 h1:Jamvg5psRIccs7FGNTlIRMkT8wgtp5eCXdBlqhYGL6U= github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= -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/config/v6 v6.0.0-beta.3 h1:G0EUzJ6Yw4UnleM6BhnOBbYPXKDHRmCJiGhC3nXDBwI= github.com/roadrunner-server/config/v6 v6.0.0-beta.3/go.mod h1:eIB+c29njpcKokXrxe483FbQOBSTNGvU3hhC6W/qYSU= github.com/roadrunner-server/endure/v2 v2.6.2 h1:sIB4kTyE7gtT3fDhuYWUYn6Vt/dcPtiA6FoNS1eS+84= github.com/roadrunner-server/endure/v2 v2.6.2/go.mod h1:t/2+xpNYgGBwhzn83y2MDhvhZ19UVq1REcvqn7j7RB8= github.com/roadrunner-server/errors v1.5.0 h1:unG7LKIZrSzkCCF3YLRLA5VyqE0KKomofXVJUXJe00g= github.com/roadrunner-server/errors v1.5.0/go.mod h1:g9fo/T2C13cWRDR9PW1r0ZAOSQfNhWAZawyfkGiaHuI= +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/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/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/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/tcplisten v1.5.2 h1:nn8yXYrhRDkfQ9AAu4V075uT4fZRmOnpxkawgE+bWPA= github.com/roadrunner-server/tcplisten v1.5.2/go.mod h1:DufGBz7Dlx2KrNe/4RukEvGMTqZKB0Uve1GztwcyyR8= github.com/rogpeppe/go-internal v1.14.1 h1:UQB4HGPB6osV0SQTLymcB4TgvyWu6ZyliaW0tI/otEQ= @@ -68,18 +56,6 @@ github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= github.com/subosito/gotenv v1.6.0 h1:9NlTDc1FTs4qu0DDq7AEtTPNw6SVm7uBMsUCUjABIf8= github.com/subosito/gotenv v1.6.0/go.mod h1:Dk4QP5c2W3ibzajGcXpNraDfq2IrhjMIvMSWPKKo0FU= -go.opentelemetry.io/auto/sdk v1.2.1 h1:jXsnJ4Lmnqd11kwkBV2LgLoFMZKizbCi5fNZ/ipaZ64= -go.opentelemetry.io/auto/sdk v1.2.1/go.mod h1:KRTj+aOaElaLi+wW1kO/DZRXwkF4C5xPbEe3ZiIhN7Y= -go.opentelemetry.io/otel v1.43.0 h1:mYIM03dnh5zfN7HautFE4ieIig9amkNANT+xcVxAj9I= -go.opentelemetry.io/otel v1.43.0/go.mod h1:JuG+u74mvjvcm8vj8pI5XiHy1zDeoCS2LB1spIq7Ay0= -go.opentelemetry.io/otel/metric v1.43.0 h1:d7638QeInOnuwOONPp4JAOGfbCEpYb+K6DVWvdxGzgM= -go.opentelemetry.io/otel/metric v1.43.0/go.mod h1:RDnPtIxvqlgO8GRW18W6Z/4P462ldprJtfxHxyKd2PY= -go.opentelemetry.io/otel/sdk v1.43.0 h1:pi5mE86i5rTeLXqoF/hhiBtUNcrAGHLKQdhg4h4V9Dg= -go.opentelemetry.io/otel/sdk v1.43.0/go.mod h1:P+IkVU3iWukmiit/Yf9AWvpyRDlUeBaRg6Y+C58QHzg= -go.opentelemetry.io/otel/sdk/metric v1.43.0 h1:S88dyqXjJkuBNLeMcVPRFXpRw2fuwdvfCGLEo89fDkw= -go.opentelemetry.io/otel/sdk/metric v1.43.0/go.mod h1:C/RJtwSEJ5hzTiUz5pXF1kILHStzb9zFlIEe85bhj6A= -go.opentelemetry.io/otel/trace v1.43.0 h1:BkNrHpup+4k4w+ZZ86CZoHHEkohws8AY+WTX09nk+3A= -go.opentelemetry.io/otel/trace v1.43.0/go.mod h1:/QJhyVBUUswCphDVxq+8mld+AvhXZLhe+8WVFxiFff0= go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto= go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE= go.uber.org/multierr v1.11.0 h1:blXXJkSxSSfBVBlC76pxqeO+LN3aDfLQo+309xJstO0= @@ -88,18 +64,10 @@ go.uber.org/zap v1.28.0 h1:IZzaP1Fv73/T/pBMLk4VutPl36uNC+OSUh3JLG3FIjo= go.uber.org/zap v1.28.0/go.mod h1:rDLpOi171uODNm/mxFcuYWxDsqWSAVkFdX4XojSKg/Q= go.yaml.in/yaml/v3 v3.0.4 h1:tfq32ie2Jv2UxXFdLJdh3jXuOzWiL1fo0bu/FbuKpbc= go.yaml.in/yaml/v3 v3.0.4/go.mod h1:DhzuOOF2ATzADvBadXxruRBLzYTpT36CKvDb3+aBEFg= -golang.org/x/net v0.56.0 h1:Rw8j/hFzGvJUZwNBXnAtf5sVDVt+65SK2C7IxCxZt5o= -golang.org/x/net v0.56.0/go.mod h1:D3Ku6r+V6JROoZK144D2XfMHFcMq/0zSfLelVTCFKec= golang.org/x/sys v0.46.0 h1:noSf2Fq6F8DBgS+LysIkx7rIExoNHJsxOAtPp4rthXw= golang.org/x/sys v0.46.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= golang.org/x/text v0.38.0 h1:sXmwo9DwP3OK9EZ7PqAdaooSGozfl/3a6/xJcbzPRhE= golang.org/x/text v0.38.0/go.mod h1:YXZt3QhHUKYT53r2lLKFIVi6Ao1jdzrTR/KQ09qyxF4= -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/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/lock_api_test.go b/tests/lock_api_test.go deleted file mode 100644 index 8db158d..0000000 --- a/tests/lock_api_test.go +++ /dev/null @@ -1,270 +0,0 @@ -package lock - -import ( - "bytes" - "context" - "crypto/tls" - "encoding/base64" - "io" - "log/slog" - "net" - "net/http" - "net/url" - "sync" - "testing" - "time" - - "connectrpc.com/connect" - lockV1 "github.com/roadrunner-server/api-go/v6/lock/v1" - "github.com/roadrunner-server/api-go/v6/lock/v1/lockV1connect" - "github.com/roadrunner-server/config/v6" - "github.com/roadrunner-server/endure/v2" - lockPlugin "github.com/roadrunner-server/lock/v6" - "github.com/roadrunner-server/logger/v6" - rpcPlugin "github.com/roadrunner-server/rpc/v6" - "github.com/stretchr/testify/require" - "golang.org/x/net/http2" - "google.golang.org/grpc" - "google.golang.org/grpc/credentials/insecure" - "google.golang.org/protobuf/encoding/protojson" - "google.golang.org/protobuf/proto" -) - -const lockAPIAddr = "127.0.0.1:6001" - -// startLockAPIContainer brings up rpc + lock + logger on lockAPIAddr. -// Returns a stop function the test must defer. -func startLockAPIContainer(t *testing.T) func() { - t.Helper() - - cont := endure.New(slog.LevelError) - cfg := &config.Plugin{ - Version: "2024.2.0", - Path: "configs/.rr-lock-api.yaml", - } - - require.NoError(t, cont.RegisterAll( - cfg, - &logger.Plugin{}, - &rpcPlugin.Plugin{}, - &lockPlugin.Plugin{}, - )) - require.NoError(t, cont.Init()) - - ch, err := cont.Serve() - require.NoError(t, err) - - wg := &sync.WaitGroup{} - stop := make(chan struct{}) - wg.Go(func() { - select { - case e := <-ch: - require.NoError(t, e.Error, "container reported error") - case <-stop: - } - }) - - time.Sleep(500 * time.Millisecond) - - return func() { - close(stop) - require.NoError(t, cont.Stop()) - wg.Wait() - } -} - -// TestLockConnectAPI exercises the lock RPCs through the Connect-RPC client -// (h2c). Go callers that import the generated lockV1connect package see -// exactly this wire shape. -func TestLockConnectAPI(t *testing.T) { - stop := startLockAPIContainer(t) - defer stop() - - httpc := &http.Client{ - Transport: &http2.Transport{ - AllowHTTP: true, - DialTLSContext: func(ctx context.Context, network, addr string, _ *tls.Config) (net.Conn, error) { - return (&net.Dialer{Timeout: 30 * time.Second}).DialContext(ctx, network, addr) - }, - }, - } - client := lockV1connect.NewLockServiceClient(httpc, "http://"+lockAPIAddr) - ctx := t.Context() - - const ( - resource = "connect-resource" - id = "connect-id" - ) - ttl := int64(30 * time.Second / time.Microsecond) - wait := int64(time.Second / time.Microsecond) - - resp, err := client.Lock(ctx, connect.NewRequest(&lockV1.LockRequest{ - Resource: resource, - Id: id, - Ttl: &ttl, - Wait: &wait, - })) - require.NoError(t, err) - require.True(t, resp.Msg.GetOk()) - - resp, err = client.Exists(ctx, connect.NewRequest(&lockV1.LockRequest{Resource: resource, Id: id})) - require.NoError(t, err) - require.True(t, resp.Msg.GetOk()) - - resp, err = client.Release(ctx, connect.NewRequest(&lockV1.LockRequest{Resource: resource, Id: id})) - require.NoError(t, err) - require.True(t, resp.Msg.GetOk()) - - resp, err = client.Exists(ctx, connect.NewRequest(&lockV1.LockRequest{Resource: resource, Id: id})) - require.NoError(t, err) - require.False(t, resp.Msg.GetOk()) -} - -// TestLockHTTPApi exercises the lock RPCs through plain HTTP/1.1 with a -// protojson body — the wire shape PHP clients use via Guzzle/curl -// (PHP has no Connect SDK). -func TestLockHTTPApi(t *testing.T) { - stop := startLockAPIContainer(t) - defer stop() - - httpc := &http.Client{Timeout: 30 * time.Second} - ctx := t.Context() - - call := func(method string, in proto.Message, out proto.Message) { - body, err := protojson.Marshal(in) - require.NoError(t, err) - - req, err := http.NewRequestWithContext(ctx, http.MethodPost, - "http://"+lockAPIAddr+"/lock.v1.LockService/"+method, bytes.NewReader(body)) - require.NoError(t, err) - req.Header.Set("Content-Type", "application/json") - - resp, err := httpc.Do(req) - require.NoError(t, err) - defer func() { _ = resp.Body.Close() }() - - respBody, err := io.ReadAll(resp.Body) - require.NoError(t, err) - require.Equalf(t, http.StatusOK, resp.StatusCode, "method=%s body=%s", method, respBody) - require.NoError(t, protojson.Unmarshal(respBody, out)) - } - - const ( - resource = "http-resource" - id = "http-id" - ) - ttl := int64(30 * time.Second / time.Microsecond) - wait := int64(time.Second / time.Microsecond) - - var lockResp lockV1.LockResponse - call("Lock", &lockV1.LockRequest{ - Resource: resource, - Id: id, - Ttl: &ttl, - Wait: &wait, - }, &lockResp) - require.True(t, lockResp.GetOk()) - - var existsResp lockV1.LockResponse - call("Exists", &lockV1.LockRequest{Resource: resource, Id: id}, &existsResp) - require.True(t, existsResp.GetOk()) - - var relResp lockV1.LockResponse - call("Release", &lockV1.LockRequest{Resource: resource, Id: id}, &relResp) - require.True(t, relResp.GetOk()) - - var existsResp2 lockV1.LockResponse - call("Exists", &lockV1.LockRequest{Resource: resource, Id: id}, &existsResp2) - require.False(t, existsResp2.GetOk()) -} - -// TestLockHTTPGetIdempotency verifies which methods accept HTTP GET. Only -// Exists is marked `option idempotency_level = NO_SIDE_EFFECTS;` in the proto, -// so Connect generates a handler that accepts GET for it. Mutating methods -// stay POST-only, so GET against them returns 405 Method Not Allowed. -func TestLockHTTPGetIdempotency(t *testing.T) { - stop := startLockAPIContainer(t) - defer stop() - - body, err := protojson.Marshal(&lockV1.LockRequest{Resource: "probe", Id: "probe"}) - require.NoError(t, err) - - q := url.Values{} - q.Set("encoding", "json") - q.Set("base64", "1") - q.Set("message", base64.URLEncoding.EncodeToString(body)) - - cases := []struct { - method string - wantStatus int - }{ - {"Exists", http.StatusOK}, - {"Lock", http.StatusMethodNotAllowed}, - {"LockRead", http.StatusMethodNotAllowed}, - {"Release", http.StatusMethodNotAllowed}, - {"ForceRelease", http.StatusMethodNotAllowed}, - {"UpdateTTL", http.StatusMethodNotAllowed}, - } - - httpc := &http.Client{Timeout: 30 * time.Second} - for _, c := range cases { - t.Run(c.method, func(t *testing.T) { - req, err := http.NewRequestWithContext(t.Context(), http.MethodGet, - "http://"+lockAPIAddr+"/lock.v1.LockService/"+c.method+"?"+q.Encode(), nil) - require.NoError(t, err) - - resp, err := httpc.Do(req) - require.NoError(t, err) - defer func() { _ = resp.Body.Close() }() - - respBody, err := io.ReadAll(resp.Body) - require.NoError(t, err) - require.Equalf(t, c.wantStatus, resp.StatusCode, - "%s via GET -> %s\n%s", c.method, resp.Status, respBody) - }) - } -} - -// TestLockGRPCApi exercises the lock RPCs through a regular gRPC client -// (google.golang.org/grpc). The same Connect handler serves gRPC framing -// off the same port — used by PHP's gRPC extension. -func TestLockGRPCApi(t *testing.T) { - stop := startLockAPIContainer(t) - defer stop() - - conn, err := grpc.NewClient(lockAPIAddr, grpc.WithTransportCredentials(insecure.NewCredentials())) - require.NoError(t, err) - defer func() { _ = conn.Close() }() - - client := lockV1.NewLockServiceClient(conn) - ctx, cancel := context.WithTimeout(t.Context(), 30*time.Second) - defer cancel() - - const ( - resource = "grpc-resource" - id = "grpc-id" - ) - ttl := int64(30 * time.Second / time.Microsecond) - wait := int64(time.Second / time.Microsecond) - - lockResp, err := client.Lock(ctx, &lockV1.LockRequest{ - Resource: resource, - Id: id, - Ttl: &ttl, - Wait: &wait, - }) - require.NoError(t, err) - require.True(t, lockResp.GetOk()) - - existsResp, err := client.Exists(ctx, &lockV1.LockRequest{Resource: resource, Id: id}) - require.NoError(t, err) - require.True(t, existsResp.GetOk()) - - relResp, err := client.Release(ctx, &lockV1.LockRequest{Resource: resource, Id: id}) - require.NoError(t, err) - require.True(t, relResp.GetOk()) - - existsResp, err = client.Exists(ctx, &lockV1.LockRequest{Resource: resource, Id: id}) - require.NoError(t, err) - require.False(t, existsResp.GetOk()) -} diff --git a/tests/rpc.go b/tests/rpc.go index da7f515..aca7801 100644 --- a/tests/rpc.go +++ b/tests/rpc.go @@ -2,114 +2,79 @@ package lock import ( "context" - "crypto/tls" "net" - "net/http" + "net/rpc" - "connectrpc.com/connect" lockV1 "github.com/roadrunner-server/api-go/v6/lock/v1" - "github.com/roadrunner-server/api-go/v6/lock/v1/lockV1connect" - "golang.org/x/net/http2" + goridgeRpc "github.com/roadrunner-server/goridge/v4/pkg/rpc" ) const lockRPCAddr = "127.0.0.1:6001" -func newLockClient() lockV1connect.LockServiceClient { - 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) - }, - }, +func newLockClient() (*rpc.Client, error) { + conn, err := new(net.Dialer).DialContext(context.Background(), "tcp", lockRPCAddr) + if err != nil { + return nil, err } - return lockV1connect.NewLockServiceClient(httpc, "http://"+lockRPCAddr) + return rpc.NewClientWithCodec(goridgeRpc.NewClientCodec(conn)), nil } -func lock(resource, id string, ttl, wait int) (bool, error) { - resp, err := newLockClient().Lock( - context.Background(), - connect.NewRequest(&lockV1.LockRequest{ - Resource: resource, - Id: id, - Ttl: new(int64(ttl)), - Wait: new(int64(wait)), - }), - ) +func call(method string, in *lockV1.LockRequest) (bool, error) { + cl, err := newLockClient() if err != nil { return false, err } - return resp.Msg.GetOk(), nil -} + defer func() { _ = cl.Close() }() -func lockRead(resource, id string, ttl, wait int) (bool, error) { - resp, err := newLockClient().LockRead( - context.Background(), - connect.NewRequest(&lockV1.LockRequest{ - Resource: resource, - Id: id, - Ttl: new(int64(ttl)), - Wait: new(int64(wait)), - }), - ) - if err != nil { + out := &lockV1.LockResponse{} + if err := cl.Call(method, in, out); err != nil { return false, err } - return resp.Msg.GetOk(), nil + return out.GetOk(), nil +} + +func lock(resource, id string, ttl, wait int) (bool, error) { + return call("lock.Lock", &lockV1.LockRequest{ + Resource: resource, + Id: id, + Ttl: new(int64(ttl)), + Wait: new(int64(wait)), + }) +} + +func lockRead(resource, id string, ttl, wait int) (bool, error) { + return call("lock.LockRead", &lockV1.LockRequest{ + Resource: resource, + Id: id, + Ttl: new(int64(ttl)), + Wait: new(int64(wait)), + }) } func release(resource, id string) (bool, error) { - resp, err := newLockClient().Release( - context.Background(), - connect.NewRequest(&lockV1.LockRequest{ - Resource: resource, - Id: id, - }), - ) - if err != nil { - return false, err - } - return resp.Msg.GetOk(), nil + return call("lock.Release", &lockV1.LockRequest{ + Resource: resource, + Id: id, + }) } func updateTTL(resource, id string, ttl int) (bool, error) { - resp, err := newLockClient().UpdateTTL( - context.Background(), - connect.NewRequest(&lockV1.LockRequest{ - Resource: resource, - Id: id, - Ttl: new(int64(ttl)), - }), - ) - if err != nil { - return false, err - } - return resp.Msg.GetOk(), nil + return call("lock.UpdateTTL", &lockV1.LockRequest{ + Resource: resource, + Id: id, + Ttl: new(int64(ttl)), + }) } func forceRelease(resource string) (bool, error) { - resp, err := newLockClient().ForceRelease( - context.Background(), - connect.NewRequest(&lockV1.LockRequest{ - Resource: resource, - }), - ) - if err != nil { - return false, err - } - return resp.Msg.GetOk(), nil + return call("lock.ForceRelease", &lockV1.LockRequest{ + Resource: resource, + }) } func exists(resource, id string) (bool, error) { - resp, err := newLockClient().Exists( - context.Background(), - connect.NewRequest(&lockV1.LockRequest{ - Resource: resource, - Id: id, - }), - ) - if err != nil { - return false, err - } - return resp.Msg.GetOk(), nil + return call("lock.Exists", &lockV1.LockRequest{ + Resource: resource, + Id: id, + }) }