From a4247b0097a5f6374054c84a68d846482a445a3c Mon Sep 17 00:00:00 2001 From: Pavel Sviderski Date: Wed, 4 Dec 2024 18:23:51 +1000 Subject: [PATCH] refactor grpc proxy backends to use machine management IP in metadata, including local responses --- internal/machine/api/proxy/director.go | 26 +++++++++++++------------- internal/machine/api/proxy/local.go | 8 +++++--- internal/machine/api/proxy/remote.go | 19 +++++++------------ 3 files changed, 25 insertions(+), 28 deletions(-) diff --git a/internal/machine/api/proxy/director.go b/internal/machine/api/proxy/director.go index a716c71a..b4a52d93 100644 --- a/internal/machine/api/proxy/director.go +++ b/internal/machine/api/proxy/director.go @@ -6,24 +6,22 @@ import ( "google.golang.org/grpc/codes" "google.golang.org/grpc/metadata" "google.golang.org/grpc/status" - "net" - "strconv" "sync" ) // Director manages routing of gRPC requests between local and remote backends. type Director struct { localBackend *LocalBackend - remotePort int + remotePort uint16 remoteBackends sync.Map // mu synchronizes access to localAddress. mu sync.RWMutex localAddress string } -func NewDirector(localSockPath string, remotePort int) *Director { +func NewDirector(localSockPath string, remotePort uint16) *Director { return &Director{ - localBackend: NewLocalBackend(localSockPath), + localBackend: NewLocalBackend(localSockPath, ""), remotePort: remotePort, } } @@ -35,6 +33,8 @@ func (d *Director) UpdateLocalAddress(addr string) { defer d.mu.Unlock() d.localAddress = addr + // Replace the local backend with the one that has local address set. + d.localBackend = NewLocalBackend(d.localBackend.sockPath, addr) } // Director implements proxy.StreamDirector for grpc-proxy, routing requests to local or remote backends based @@ -60,17 +60,17 @@ func (d *Director) Director(ctx context.Context, fullMethodName string) (proxy.M d.mu.RLock() localAddress := d.localAddress + localBackend := d.localBackend d.mu.RUnlock() backends := make([]proxy.Backend, len(machines)) for i, addr := range machines { if addr == localAddress { - backends[i] = d.localBackend + backends[i] = localBackend continue } - target := net.JoinHostPort(addr, strconv.Itoa(d.remotePort)) - backend, err := d.remoteBackend(target) + backend, err := d.remoteBackend(addr) if err != nil { return proxy.One2One, nil, status.Error(codes.Internal, err.Error()) } @@ -83,18 +83,18 @@ func (d *Director) Director(ctx context.Context, fullMethodName string) (proxy.M return proxy.One2Many, backends, nil } -// remoteBackend returns a RemoteBackend for the given target from the cache or creates a new one. -func (d *Director) remoteBackend(target string) (*RemoteBackend, error) { - b, ok := d.remoteBackends.Load(target) +// remoteBackend returns a RemoteBackend for the given address from the cache or creates a new one. +func (d *Director) remoteBackend(addr string) (*RemoteBackend, error) { + b, ok := d.remoteBackends.Load(addr) if ok { return b.(*RemoteBackend), nil } - backend, err := NewRemoteBackend(target) + backend, err := NewRemoteBackend(addr, d.remotePort) if err != nil { return nil, err } - existing, loaded := d.remoteBackends.LoadOrStore(target, backend) + existing, loaded := d.remoteBackends.LoadOrStore(addr, backend) if loaded { // A concurrent remoteBackend call built a different backend. backend.Close() diff --git a/internal/machine/api/proxy/local.go b/internal/machine/api/proxy/local.go index d5638d6c..dee61744 100644 --- a/internal/machine/api/proxy/local.go +++ b/internal/machine/api/proxy/local.go @@ -20,11 +20,13 @@ type LocalBackend struct { var _ proxy.Backend = (*LocalBackend)(nil) -// NewLocalBackend returns a new LocalBackend for the given Unix socket path. -func NewLocalBackend(sockPath string) *LocalBackend { +// NewLocalBackend returns a new LocalBackend for the given Unix socket path. The addr parameter is the local address +// of the current machine which could be empty if it's not known. The address is used to populate response metadata +// in one2many mode. +func NewLocalBackend(sockPath, addr string) *LocalBackend { return &LocalBackend{ One2ManyResponder: One2ManyResponder{ - machine: "local", + machine: addr, }, sockPath: sockPath, } diff --git a/internal/machine/api/proxy/remote.go b/internal/machine/api/proxy/remote.go index 0b47fbb0..8db7983f 100644 --- a/internal/machine/api/proxy/remote.go +++ b/internal/machine/api/proxy/remote.go @@ -8,7 +8,6 @@ import ( "google.golang.org/grpc/backoff" "google.golang.org/grpc/credentials/insecure" "google.golang.org/grpc/metadata" - "net" "net/netip" "sync" "time" @@ -29,22 +28,18 @@ type RemoteBackend struct { var _ proxy.Backend = (*RemoteBackend)(nil) -// NewRemoteBackend creates a new instance of RemoteBackend for the given target which must have the format [IPv6]:port. -func NewRemoteBackend(target string) (*RemoteBackend, error) { - host, _, err := net.SplitHostPort(target) - if err != nil { - return nil, fmt.Errorf("target must have the format [IPv6]:port: %s", target) - } - addr, err := netip.ParseAddr(host) - if err != nil || !addr.Is6() { - return nil, fmt.Errorf("target host must be a valid IPv6 address: %s", host) +// NewRemoteBackend creates a new instance of RemoteBackend for the given IPv6 address and port. +func NewRemoteBackend(addr string, port uint16) (*RemoteBackend, error) { + ip, err := netip.ParseAddr(addr) + if err != nil || !ip.Is6() { + return nil, fmt.Errorf("address must be a valid IPv6 address: %s", addr) } return &RemoteBackend{ One2ManyResponder: One2ManyResponder{ - machine: target, + machine: addr, }, - target: target, + target: netip.AddrPortFrom(ip, port).String(), }, nil }