From 71070e6e68123e5ca42b36cdd401d90c3963f3da Mon Sep 17 00:00:00 2001 From: Zexi Li Date: Fri, 20 Jun 2025 19:16:55 +0800 Subject: [PATCH] feat(region,host): support port mapping rule (#22748) --- pkg/apis/compute/guestnetwork.go | 23 +++- pkg/compute/models/networks.go | 19 ++- pkg/hostman/guestman/pod.go | 15 +++ pkg/hostman/guestman/pod_sync_loop.go | 2 +- pkg/hostman/guestman/portmapping.go | 179 ++++++++++++++++++++++++++ 5 files changed, 233 insertions(+), 5 deletions(-) diff --git a/pkg/apis/compute/guestnetwork.go b/pkg/apis/compute/guestnetwork.go index 532c763c16..2c97a9beff 100644 --- a/pkg/apis/compute/guestnetwork.go +++ b/pkg/apis/compute/guestnetwork.go @@ -180,14 +180,33 @@ type GuestPortMappingPortRange struct { End int `json:"end"` } +type GuestPortMappingEnvValueFrom string + +const ( + GuestPortMappingEnvValueFromPort GuestPortMappingEnvValueFrom = "port" + GuestPortMappingEnvValueFromHostPort GuestPortMappingEnvValueFrom = "host_port" +) + +type GuestPortMappingEnv struct { + Key string `json:"key"` + ValueFrom GuestPortMappingEnvValueFrom `json:"value_from"` +} + type GuestPortMapping struct { - Protocol GuestPortMappingProtocol `json:"protocol"` + Protocol GuestPortMappingProtocol `json:"protocol"` + // 容器内部 Port 端口范围 1-65535,-1表示由宿主机自动分配和 HostPort 相同的端口 Port int `json:"port"` HostPort *int `json:"host_port,omitempty"` HostIp string `json:"host_ip"` HostPortRange *GuestPortMappingPortRange `json:"host_port_range,omitempty"` // whitelist for remote ips - RemoteIps []string `json:"remote_ips"` + RemoteIps []string `json:"remote_ips"` + Rule *GuestPortMappingRule `json:"rule,omitempty"` + Envs []GuestPortMappingEnv `json:"envs,omitempty"` +} + +type GuestPortMappingRule struct { + FirstPortOffset *int `json:"first_port_offset"` } type GuestPortMappings []*GuestPortMapping diff --git a/pkg/compute/models/networks.go b/pkg/compute/models/networks.go index aad6782706..4a680b6801 100644 --- a/pkg/compute/models/networks.go +++ b/pkg/compute/models/networks.go @@ -1233,8 +1233,11 @@ func validatePortMapping(pm *api.GuestPortMapping) error { return errors.Wrap(err, "validate host_port") } } - if err := validatePort(pm.Port, 1, 65535); err != nil { - return errors.Wrap(err, "validate port") + // -1 端口表示自动分配 + if pm.Port != -1 { + if err := validatePort(pm.Port, 1, 65535); err != nil { + return errors.Wrap(err, "validate port") + } } if pm.Protocol == "" { pm.Protocol = api.GuestPortMappingProtocolTCP @@ -1249,6 +1252,18 @@ func validatePortMapping(pm *api.GuestPortMapping) error { } } } + if pm.Rule != nil { + if pm.Rule.FirstPortOffset != nil { + if *pm.Rule.FirstPortOffset < 0 { + return httperrors.NewInputParameterError("first port offset %d is less than 0", *pm.Rule.FirstPortOffset) + } + } + } + for _, env := range pm.Envs { + if env.ValueFrom != api.GuestPortMappingEnvValueFromPort && env.ValueFrom != api.GuestPortMappingEnvValueFromHostPort { + return httperrors.NewInputParameterError("invalid value from %s", env.ValueFrom) + } + } return nil } diff --git a/pkg/hostman/guestman/pod.go b/pkg/hostman/guestman/pod.go index 47d99d9e39..408c00aacd 100644 --- a/pkg/hostman/guestman/pod.go +++ b/pkg/hostman/guestman/pod.go @@ -1892,6 +1892,21 @@ func (s *sPodGuestInstance) createContainer(ctx context.Context, userCred mcclie Key: envKey, Value: envVal, }) + for _, pEnv := range pm.Envs { + pEnvVal := "" + switch pEnv.ValueFrom { + case computeapi.GuestPortMappingEnvValueFromHostPort: + pEnvVal = fmt.Sprintf("%d", *pm.HostPort) + case computeapi.GuestPortMappingEnvValueFromPort: + pEnvVal = fmt.Sprintf("%d", pm.Port) + default: + return "", httperrors.NewInputParameterError("invalid value from %s", pEnv.ValueFrom) + } + ctrCfg.Envs = append(ctrCfg.Envs, &runtimeapi.KeyValue{ + Key: pEnv.Key, + Value: pEnvVal, + }) + } } } if s.GetDesc().HostAccessIp != "" { diff --git a/pkg/hostman/guestman/pod_sync_loop.go b/pkg/hostman/guestman/pod_sync_loop.go index 822b7ca1c6..59c086a60a 100644 --- a/pkg/hostman/guestman/pod_sync_loop.go +++ b/pkg/hostman/guestman/pod_sync_loop.go @@ -150,7 +150,7 @@ func (m *SGuestManager) syncContainerLoopIteration(plegCh chan *pleg.PodLifecycl log.Warningf("can not find pod manager by %s", jsonutils.Marshal(e)) return } - if podMan.(*sPodGuestInstance).IsDirtyShutdown() { + if podMan.(*sPodGuestInstance).isPodDirtyShutdown() { log.Infof("pod %s is dirty shutdown, waiting it to started", podMan.GetName()) return } diff --git a/pkg/hostman/guestman/portmapping.go b/pkg/hostman/guestman/portmapping.go index d3bbfa0b95..d223f64075 100644 --- a/pkg/hostman/guestman/portmapping.go +++ b/pkg/hostman/guestman/portmapping.go @@ -156,6 +156,22 @@ func (m *portMappingManager) getOtherGuestsUsedPorts(gst GuestRuntimeInstance) ( func (m *portMappingManager) allocatePortMappings(gst GuestRuntimeInstance, input compute.GuestPortMappings) (compute.GuestPortMappings, error) { result := make([]*compute.GuestPortMapping, len(input)) allocPorts := make(map[compute.GuestPortMappingProtocol]sets.Int) + + // 检查是否有需要按规则分配的端口映射 + hasRuleMapping := false + for _, pm := range input { + if pm.Rule != nil && pm.Rule.FirstPortOffset != nil { + hasRuleMapping = true + break + } + } + + if hasRuleMapping { + // 如果有规则映射,需要先找到第一个空闲端口,然后按偏移量分配 + return m.allocatePortMappingsWithRule(gst, input, allocPorts) + } + + // 原有的分配逻辑 for idx := range input { data := input[idx] if _, ok := allocPorts[data.Protocol]; !ok { @@ -171,6 +187,163 @@ func (m *portMappingManager) allocatePortMappings(gst GuestRuntimeInstance, inpu return result, nil } +func (m *portMappingManager) allocatePortMappingsWithRule(gst GuestRuntimeInstance, input compute.GuestPortMappings, allocPorts map[compute.GuestPortMappingProtocol]sets.Int) (compute.GuestPortMappings, error) { + result := make([]*compute.GuestPortMapping, len(input)) + + // 按协议分组,分别处理 + indices := make([]*compute.GuestPortMapping, 0) + for idx, pm := range input { + if pm.Rule != nil && pm.Rule.FirstPortOffset != nil { + indices = append(indices, input[idx]) + } + } + + // 为每个协议组分配端口 + if err := m.allocateProtocolGroupWithRule(gst, input, result, indices, allocPorts); err != nil { + return nil, errors.Wrapf(err, "allocate portmappings with rule: %s", jsonutils.Marshal(indices)) + } + + // 处理没有规则的端口映射 + for idx, pm := range input { + if pm.Rule == nil || pm.Rule.FirstPortOffset == nil { + if _, ok := allocPorts[pm.Protocol]; !ok { + allocPorts[pm.Protocol] = sets.NewInt() + } + allocatedPm, err := m.allocatePortMapping(gst, pm, allocPorts) + if err != nil { + return nil, errors.Wrapf(err, "get port mapping %s", jsonutils.Marshal(pm)) + } + result[idx] = allocatedPm + allocPorts[pm.Protocol].Insert(*allocatedPm.HostPort) + } + } + + return result, nil +} + +func (m *portMappingManager) allocateProtocolGroupWithRule(gst GuestRuntimeInstance, input compute.GuestPortMappings, result compute.GuestPortMappings, indices []*compute.GuestPortMapping, allocPorts map[compute.GuestPortMappingProtocol]sets.Int) error { + // 获取其他虚拟机已使用的端口 + otherPorts, err := m.getOtherGuestsUsedPorts(gst) + if err != nil { + return errors.Wrap(err, "getOtherGuestsUsedPorts") + } + + // 获取当前协议已分配的端口 + usedPorts := map[compute.GuestPortMappingProtocol]sets.Int{ + compute.GuestPortMappingProtocolTCP: sets.NewInt(), + compute.GuestPortMappingProtocolUDP: sets.NewInt(), + } + for proto, ports := range otherPorts { + usedPorts[proto].Insert(ports.List()...) + } + for proto, allocPortsSet := range allocPorts { + if ports, ok := usedPorts[proto]; ok { + ports.Insert(allocPortsSet.List()...) + usedPorts[proto] = ports + } else { + usedPorts[proto] = sets.NewInt(allocPortsSet.List()...) + } + } + + // 确定端口范围 + start := compute.GUEST_PORT_MAPPING_RANGE_START + end := compute.GUEST_PORT_MAPPING_RANGE_END + + // 尝试不同的 basePort,直到找到满足所有规则要求的端口 + success := false + for basePort := start; basePort <= end; basePort++ { + // 检查这个 basePort 是否能满足所有规则要求 + if m.canAllocateWithBasePort(basePort, input, indices, usedPorts) { + // 分配端口 + if err := m.allocateWithBasePort(basePort, input, result, indices, usedPorts, allocPorts); err != nil { + // 如果分配失败,继续尝试下一个 basePort + continue + } + success = true + break + } + } + + if !success { + return errors.Errorf("cannot find suitable base port for protocol %s in range %d-%d", indices[0].Protocol, start, end) + } + + return nil +} + +func (m *portMappingManager) checkPortIsUsed(port int, protocol compute.GuestPortMappingProtocol, usedPorts map[compute.GuestPortMappingProtocol]sets.Int) bool { + portProtocol := getport.TCP + if protocol == compute.GuestPortMappingProtocolUDP { + portProtocol = getport.UDP + } + if _, ok := usedPorts[protocol]; !ok { + usedPorts[protocol] = sets.NewInt() + } + return usedPorts[protocol].Has(port) || getport.IsPortUsed(portProtocol, "", port) +} + +func (m *portMappingManager) canAllocateWithBasePort(basePort int, input compute.GuestPortMappings, indices []*compute.GuestPortMapping, usedPorts map[compute.GuestPortMappingProtocol]sets.Int) bool { + baseProtocol := indices[0].Protocol + + // 检查 basePort 本身是否可用 + if m.checkPortIsUsed(basePort, baseProtocol, usedPorts) { + return false + } + + // 检查所有设置了规则的端口是否都可用 + for _, pm := range indices { + offset := *pm.Rule.FirstPortOffset + targetPort := basePort + offset + + // 检查目标端口是否在范围内 + if targetPort > compute.GUEST_PORT_MAPPING_RANGE_END { + return false + } + + // 检查目标端口是否已被使用 + if m.checkPortIsUsed(targetPort, pm.Protocol, usedPorts) { + return false + } + } + + return true +} + +func (m *portMappingManager) allocateWithBasePort(basePort int, input compute.GuestPortMappings, result compute.GuestPortMappings, indices []*compute.GuestPortMapping, usedPorts, allocPorts map[compute.GuestPortMappingProtocol]sets.Int) error { + // 分配所有设置了规则的端口 + for idx, _ := range indices { + pm := input[idx] + offset := *pm.Rule.FirstPortOffset + targetPort := basePort + offset + + // 再次检查端口可用性(双重检查) + if m.checkPortIsUsed(targetPort, pm.Protocol, usedPorts) { + return errors.Errorf("port %d is not available for protocol %s", targetPort, pm.Protocol) + } + + // 创建分配的端口映射 + runtimePm := &compute.GuestPortMapping{} + if err := jsonutils.Marshal(pm).Unmarshal(runtimePm); err != nil { + return errors.Wrap(err, "unmarshal to runtime port mapping") + } + + runtimePm.HostPort = &targetPort + if runtimePm.Port == -1 { + runtimePm.Port = targetPort + } + result[idx] = runtimePm + + // 更新已使用端口集合 + usedPorts[pm.Protocol].Insert(targetPort) + if _, ok := allocPorts[pm.Protocol]; !ok { + allocPorts[pm.Protocol] = sets.NewInt() + } + allocPorts[pm.Protocol].Insert(targetPort) + } + + return nil +} + func (m *portMappingManager) allocatePortMapping(gst GuestRuntimeInstance, pm *compute.GuestPortMapping, allocPorts map[compute.GuestPortMappingProtocol]sets.Int) (*compute.GuestPortMapping, error) { otherPorts, err := m.getOtherGuestsUsedPorts(gst) if err != nil { @@ -210,6 +383,9 @@ func (m *portMappingManager) allocatePortMapping(gst GuestRuntimeInstance, pm *c return nil, errors.Errorf("%s host_port %d is already allocated", pm.Protocol, *pm.HostPort) } } + if runtimePm.Port == -1 { + runtimePm.Port = *pm.HostPort + } return runtimePm, nil } else { start := compute.GUEST_PORT_MAPPING_RANGE_START @@ -231,6 +407,9 @@ func (m *portMappingManager) allocatePortMapping(gst GuestRuntimeInstance, pm *c return nil, errors.Wrapf(err, "listen %s port inside %d and %d", pm.Protocol, start, end) } runtimePm.HostPort = &portResult.Port + if runtimePm.Port == -1 { + runtimePm.Port = portResult.Port + } return runtimePm, nil } }