feat(region,host): support port mapping rule (#22748)

This commit is contained in:
Zexi Li
2025-06-20 19:16:55 +08:00
committed by GitHub
parent f50c4a5c78
commit 71070e6e68
5 changed files with 233 additions and 5 deletions
+21 -2
View File
@@ -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
+17 -2
View File
@@ -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
}
+15
View File
@@ -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 != "" {
+1 -1
View File
@@ -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
}
+179
View File
@@ -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
}
}