mirror of
https://github.com/psviderski/uncloud.git
synced 2026-08-26 11:03:34 +00:00
refactor grpc proxy backends to use machine management IP in metadata, including local responses
This commit is contained in:
@@ -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()
|
||||
|
||||
@@ -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,
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user