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
279 changes: 278 additions & 1 deletion pkg/agent/server/guest.go
Original file line number Diff line number Diff line change
Expand Up @@ -21,9 +21,13 @@ import (

"github.com/digitalocean/go-openvswitch/ovs"

"yunion.io/x/jsonutils"
"yunion.io/x/log"
"yunion.io/x/pkg/errors"
"yunion.io/x/sdnagent/pkg/agent/utils"

computeapi "yunion.io/x/onecloud/pkg/apis/compute"
fwdpb "yunion.io/x/onecloud/pkg/hostman/guestman/forwarder/api"
"yunion.io/x/onecloud/pkg/mcclient/auth"
mcclient_modules "yunion.io/x/onecloud/pkg/mcclient/modules/compute"
)
Expand All @@ -38,6 +42,10 @@ type Guest struct {
*utils.Guest
watcher *serversWatcher
lastSeenPending *time.Time

// vpcPortMapFwds tracks persistent ovnMd forwards opened for VPC nic port mappings.
// key: proto/bindPort/remoteAddr/remotePort
vpcPortMapFwds map[string]*fwdpb.OpenResponse
}

func NewGuest(guest *utils.Guest, watcher *serversWatcher) *Guest {
Expand Down Expand Up @@ -176,6 +184,12 @@ func (g *Guest) updateClassicFlows(ctx context.Context) (err error) {
flowman.updateFlows(ctx, g.Who(), flows)
}
}
if err2 := utils.SyncGuestPortMappingDNAT(g.Id, g.NICs); err2 != nil {
log.Errorf("guest %s sync port mapping DNAT: %v", g.Id, err2)
if err == nil {
err = err2
}
}
return
}

Expand All @@ -191,6 +205,9 @@ func (g *Guest) clearClassicFlows(ctx context.Context) {
flowman.updateFlows(ctx, g.Who(), []*ovs.Flow{})
}
}
if err := utils.ClearGuestPortMappingDNAT(g.Id); err != nil {
log.Errorf("guest %s clear port mapping DNAT: %v", g.Id, err)
}
g.clearPending()
}

Expand All @@ -217,6 +234,16 @@ func (g *Guest) updateOvn(ctx context.Context) {
if g.HostConfig.DisableLocalVpc {
return
}
if g.watcher.ovnMan == nil || g.watcher.ovnMdMan == nil {
return
}

if len(g.VpcNICs) == 0 {
g.clearVpcPortMappingForwards(ctx)
g.watcher.ovnMan.SetGuestNICs(ctx, g.Id, nil)
g.watcher.ovnMdMan.SetGuestNICs(ctx, g.Id, nil)
return
}

if len(g.VpcNICs) > 0 && g.HostId != "" {
ovnMan := g.watcher.ovnMan
Expand All @@ -225,6 +252,11 @@ func (g *Guest) updateOvn(ctx context.Context) {

ovnMdMan := g.watcher.ovnMdMan
ovnMdMan.SetGuestNICs(ctx, g.Id, g.VpcNICs)

if err := g.syncVpcPortMappingForwards(ctx); err != nil {
log.Errorf("guest %s sync vpc port mapping forwards: %v", g.Id, err)
g.setPending()
}
}
}

Expand All @@ -233,13 +265,245 @@ func (g *Guest) clearOvn(ctx context.Context) {
return
}

g.clearVpcPortMappingForwards(ctx)

ovnMan := g.watcher.ovnMan
ovnMan.SetGuestNICs(ctx, g.Id, nil)

ovnMdMan := g.watcher.ovnMdMan
ovnMdMan.SetGuestNICs(ctx, g.Id, nil)
}

func (g *Guest) syncVpcPortMappingForwards(ctx context.Context) error {
if g.watcher.ovnMdMan == nil {
return nil
}
if g.vpcPortMapFwds == nil {
g.vpcPortMapFwds = map[string]*fwdpb.OpenResponse{}
}
var errs []error
aliveIPs := map[string]bool{}
for _, nic := range g.VpcNICs {
if nic.IP != "" {
aliveIPs[nic.IP] = true
}
if err := g.syncNicPortMappingForwards(ctx, nic); err != nil {
log.Errorf("guest %s nic %s sync vpc port mapping forwards: %v", g.Id, nic.NetId, err)
errs = append(errs, err)
}
}
// close forwards for VPC nics that are gone
for key, fwd := range g.vpcPortMapFwds {
if aliveIPs[fwd.RemoteAddr] {
continue
}
if err := g.closePortMappingForward(ctx, g.vpcPortMapFwdNetId(fwd), fwd); err != nil {
log.Errorf("guest %s close stale port mapping forward %s:%d -> %s:%d: %v",
g.Id, fwd.BindAddr, fwd.BindPort, fwd.RemoteAddr, fwd.RemotePort, err)
}
delete(g.vpcPortMapFwds, key)
}
if len(errs) > 0 {
return errors.NewAggregate(errs)
}
return nil
}

func (g *Guest) clearVpcPortMappingForwards(ctx context.Context) {
if g.watcher.ovnMdMan == nil {
return
}
for key, fwd := range g.vpcPortMapFwds {
netId := g.vpcPortMapFwdNetId(fwd)
if err := g.closePortMappingForward(ctx, netId, fwd); err != nil {
log.Errorf("guest %s clear port mapping forward %s:%d -> %s:%d: %v",
g.Id, fwd.BindAddr, fwd.BindPort, fwd.RemoteAddr, fwd.RemotePort, err)
}
delete(g.vpcPortMapFwds, key)
}
}

func (g *Guest) vpcPortMapFwdNetId(fwd *fwdpb.OpenResponse) string {
if fwd == nil {
return ""
}
if fwd.NetId != "" {
return fwd.NetId
}
for _, nic := range g.VpcNICs {
if nic.IP == fwd.RemoteAddr {
return nic.NetId
}
}
return ""
}

func (g *Guest) syncNicPortMappingForwards(ctx context.Context, nic *utils.GuestNIC) error {
if nic == nil || nic.NetId == "" || nic.IP == "" {
return nil
}

desired := map[string]*fwdpb.OpenRequest{}
for _, pm := range nic.PortMappings {
req, key, ok := portMappingToOpenRequest(nic, pm)
if !ok {
continue
}
req.BindAddr = g.resolvePortMappingBindAddr(req.BindAddr)
// key is independent of bind addr
desired[key] = req
log.Infof("guest %s port mapping forward %s", key, jsonutils.Marshal(req).String())
}

// Close forwards we previously opened for this nic but are no longer desired.
// Do not close untracked forwards (e.g. ephemeral ssh) returned by ListByRemote.
for key, fwd := range g.vpcPortMapFwds {
if fwd.RemoteAddr != nic.IP {
continue
}
if _, ok := desired[key]; ok {
continue
}
if err := g.closePortMappingForward(ctx, nic.NetId, fwd); err != nil {
log.Errorf("guest %s close port mapping forward %s:%d -> %s:%d: %v",
g.Id, fwd.BindAddr, fwd.BindPort, fwd.RemoteAddr, fwd.RemotePort, err)
}
delete(g.vpcPortMapFwds, key)
}

existingByKey := map[string]*fwdpb.OpenResponse{}
existing, err := g.listNicPortMappingForwards(ctx, nic)
if err != nil {
// md server may not be ready yet; still try Open below
log.Warningf("guest %s nic %s list port mapping forwards: %v", g.Id, nic.NetId, err)
} else {
for _, fwd := range existing {
key := portMappingFwdKey(fwd.Proto, fwd.BindPort, fwd.RemoteAddr, fwd.RemotePort)
existingByKey[key] = fwd
}
}

var openErrs []error
for key, req := range desired {
if fwd, ok := existingByKey[key]; ok {
if fwd.NetId == "" {
fwd.NetId = nic.NetId
}
g.vpcPortMapFwds[key] = fwd
continue
}
// tracked or not, forward is not running — (re)open
delete(g.vpcPortMapFwds, key)
pbresp, err := g.watcher.ovnMdMan.ForwardRequest(ctx, ovnMdFwdReq{pbreq: req})
if err != nil {
log.Errorf("guest %s open port mapping forward %s:%d -> %s:%d: %v",
g.Id, req.BindAddr, req.BindPort, req.RemoteAddr, req.RemotePort, err)
openErrs = append(openErrs, err)
continue
}
fwd, ok := pbresp.(*fwdpb.OpenResponse)
if !ok || fwd == nil {
continue
}
if fwd.NetId == "" {
fwd.NetId = nic.NetId
}
g.vpcPortMapFwds[key] = fwd
log.Infof("guest %s port mapping forward ready %s %s:%d -> %s:%d",
g.Id, fwd.Proto, fwd.BindAddr, fwd.BindPort, fwd.RemoteAddr, fwd.RemotePort)
}
if len(openErrs) > 0 {
return errors.NewAggregate(openErrs)
}
return nil
}

func (g *Guest) listNicPortMappingForwards(ctx context.Context, nic *utils.GuestNIC) ([]*fwdpb.OpenResponse, error) {
if nic == nil || nic.NetId == "" || nic.IP == "" {
return nil, nil
}
pbresp, err := g.watcher.ovnMdMan.ForwardRequest(ctx, ovnMdFwdReq{
pbreq: &fwdpb.ListByRemoteRequest{
NetId: nic.NetId,
RemoteAddr: nic.IP,
},
})
if err != nil {
return nil, err
}
listResp, ok := pbresp.(*fwdpb.ListByRemoteResponse)
if !ok || listResp == nil {
return nil, nil
}
return listResp.Forwards, nil
}

func (g *Guest) closePortMappingForward(ctx context.Context, netId string, fwd *fwdpb.OpenResponse) error {
if fwd == nil {
return nil
}
if netId == "" {
netId = fwd.NetId
}
if netId == "" {
return fmt.Errorf("missing netId for forward %s:%d", fwd.BindAddr, fwd.BindPort)
}
_, err := g.watcher.ovnMdMan.ForwardRequest(ctx, ovnMdFwdReq{
pbreq: &fwdpb.CloseRequest{
NetId: netId,
Proto: fwd.Proto,
BindAddr: fwd.BindAddr,
BindPort: fwd.BindPort,
},
})
return err
}

func portMappingToOpenRequest(nic *utils.GuestNIC, pm *computeapi.GuestPortMapping) (*fwdpb.OpenRequest, string, bool) {
if nic == nil || pm == nil || pm.HostPort == nil || nic.IP == "" || pm.Port <= 0 || nic.NetId == "" {
return nil, "", false
}
proto := string(pm.Protocol)
if proto == "" {
proto = string(computeapi.GuestPortMappingProtocolTCP)
}
// ovnMdForward currently only supports tcp
if proto != string(computeapi.GuestPortMappingProtocolTCP) {
log.Warningf("skip vpc port mapping forward: unsupported proto %s nic=%s port=%d", proto, nic.IP, pm.Port)
return nil, "", false
}
bindAddr := pm.HostIp
if bindAddr == "" || bindAddr == "0.0.0.0" {
// follow hostman OpenForward: bind master IP for reachable host-side proxy
bindAddr = ""
}
req := &fwdpb.OpenRequest{
NetId: nic.NetId,
Proto: proto,
BindAddr: bindAddr,
BindPort: uint32(*pm.HostPort),
RemoteAddr: nic.IP,
RemotePort: uint32(pm.Port),
}
return req, portMappingFwdKey(req.Proto, req.BindPort, req.RemoteAddr, req.RemotePort), true
}

func (g *Guest) resolvePortMappingBindAddr(bindAddr string) string {
if bindAddr != "" && bindAddr != "0.0.0.0" {
return bindAddr
}
if g.HostConfig != nil {
if nic := g.HostConfig.MasterNic(); nic != nil && nic.Addr != "" {
return nic.Addr
}
}
return "0.0.0.0"
}

func portMappingFwdKey(proto string, bindPort uint32, remoteAddr string, remotePort uint32) string {
return fmt.Sprintf("%s/%d/%s/%d", proto, bindPort, remoteAddr, remotePort)
}

func (g *Guest) UpdateSettings(ctx context.Context, sync bool) {
start := time.Now()
err := g.refresh(ctx)
Expand All @@ -255,7 +519,20 @@ func (g *Guest) UpdateSettings(ctx context.Context, sync bool) {
if g.HostId != "" {
g.watcher.agent.HostId(g.HostId)
}
case errNotRunning, errPortNotReady, errVolatileHost:
case errPortNotReady:
// Classic/OVS port may be late, but VPC metadata forward does not need PortNo.
// Keep OVN + port-mapping forwards for running VPC guests instead of clearing them.
if g.Running() && len(g.VpcNICs) > 0 && !g.HostConfig.DisableLocalVpc {
log.Debugf("guest %s(%s) classic port not ready, still update OVN/port mapping", g.Name, g.Id)
g.clearClassicFlows(ctx)
g.clearTc(ctx)
g.updateOvn(ctx)
g.setPending()
break
}
log.Debugf("guest %s(%s) ClearSettings due to g.refresh %s", g.Name, g.Id, err)
g.ClearSettings(ctx)
case errNotRunning, errVolatileHost:
log.Debugf("guest %s(%s) ClearSettings due to g.refresh %s", g.Name, g.Id, err)
g.ClearSettings(ctx)
default:
Expand Down
Loading
Loading