Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 2 additions & 12 deletions go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -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
40 changes: 2 additions & 38 deletions go.sum
Original file line number Diff line number Diff line change
@@ -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=
7 changes: 2 additions & 5 deletions plugin.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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}
}
87 changes: 40 additions & 47 deletions rpc.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,6 @@ import (
"errors"
"time"

"connectrpc.com/connect"
lockV1 "github.com/roadrunner-server/api-go/v6/lock/v1"
)

Expand All @@ -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
}
9 changes: 0 additions & 9 deletions tests/configs/.rr-lock-api.yaml

This file was deleted.

12 changes: 4 additions & 8 deletions tests/go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
)
Loading
Loading