Compare commits
13
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
1a26a28876 | ||
|
|
d2b1a21a07 | ||
|
|
8718d27fc7 | ||
|
|
f991c2fe15 | ||
|
|
cf30022133 | ||
|
|
1ce335d3db | ||
|
|
10d741ceb6 | ||
|
|
23ee9db817 | ||
|
|
e6054bc8ea | ||
|
|
3935f5e2e8 | ||
|
|
1a5d6c8b37 | ||
|
|
c40db0d85f | ||
|
|
6ba3ad4281 |
@@ -1,4 +0,0 @@
|
||||
owner: KubelanCloud
|
||||
git-repo: kks-provider-plugin
|
||||
pages-branch: gh-pages
|
||||
pages-index-path: index.yaml
|
||||
@@ -20,8 +20,7 @@ jobs:
|
||||
|
||||
- name: Container image name
|
||||
run: |
|
||||
owner=$(echo "${GITHUB_REPOSITORY_OWNER}" | tr '[:upper:]' '[:lower:]')
|
||||
echo "IMAGE_NAME=${REGISTRY_HOST}/kloude/kks-provider-plugin" >> "$GITHUB_ENV"
|
||||
echo "IMAGE_NAME=${REGISTRY_HOST}/kloude-public/kks-provider-plugin" >> "$GITHUB_ENV"
|
||||
echo "APP_VERSION=$(grep '^appVersion:' charts/kks-provider-plugin/Chart.yaml | awk '{print $2}' | tr -d '\"')" >> "$GITHUB_ENV"
|
||||
|
||||
- uses: docker/setup-buildx-action@v3
|
||||
@@ -52,11 +51,3 @@ jobs:
|
||||
push: ${{ github.event_name != 'pull_request' }}
|
||||
tags: ${{ steps.meta.outputs.tags }}
|
||||
labels: ${{ steps.meta.outputs.labels }}
|
||||
|
||||
- name: Make container image public
|
||||
if: github.event_name != 'pull_request'
|
||||
env:
|
||||
GH_TOKEN: ${{ secrets.GITHUB_TOKEN }}
|
||||
run: |
|
||||
owner=$(echo "${GITHUB_REPOSITORY_OWNER}" | tr '[:upper:]' '[:lower:]')
|
||||
gh api --method PATCH "/orgs/${owner}/packages/container/kks-provider-plugin/visibility" -f visibility=public || true
|
||||
|
||||
@@ -96,27 +96,24 @@ jobs:
|
||||
"${api}/releases/${release_id}/assets?name=${file_name}"
|
||||
done
|
||||
|
||||
- name: Log in to registry
|
||||
uses: docker/login-action@v3
|
||||
with:
|
||||
registry: registry.kloude.ir
|
||||
username: ${{ github.actor }}
|
||||
password: ${{ secrets.GITHUB_TOKEN }}
|
||||
|
||||
- name: Push chart to OCI registry
|
||||
run: |
|
||||
set -euo pipefail
|
||||
|
||||
shopt -s nullglob
|
||||
charts=(.cr-release-packages/*.tgz)
|
||||
if [ ${#charts[@]} -eq 0 ]; then
|
||||
echo "No packaged charts to push"
|
||||
exit 0
|
||||
fi
|
||||
for chart in "${charts[@]}"; do
|
||||
helm push "$chart" "oci://registry.kloude.ir/kloude/kks-provider-plugin-charts"
|
||||
done
|
||||
|
||||
- name: Make OCI chart public
|
||||
env:
|
||||
GH_TOKEN: ${{ secrets.GITHUB_TOKEN }}
|
||||
run: |
|
||||
gh api --method PATCH "/orgs/kloude/kks-provider-plugin/packages/container/charts%2Fkks-provider-plugin/visibility" -f visibility=public || true
|
||||
repo="oci://registry.kloude.ir/kloude-public/kks-provider-plugin-charts"
|
||||
|
||||
for chart in "${charts[@]}"; do
|
||||
if ! helm push --plain-http \
|
||||
--username "${{ secrets.REGISTRY_USER }}" \
|
||||
--password "${{ secrets.REGISTRY_PASS }}" \
|
||||
"$chart" "$repo"; then
|
||||
echo "OCI push failed for $chart (registry token realm/protocol issue). Continuing because chart artifact is already published to Gitea release assets."
|
||||
fi
|
||||
done
|
||||
|
||||
@@ -48,52 +48,8 @@ jobs:
|
||||
helm package "${chart_dir}" --destination dist
|
||||
echo "CHART_PACKAGE=dist/kks-provider-plugin-${version}.tgz" >> "$GITHUB_ENV"
|
||||
|
||||
- name: Build changelog
|
||||
id: changelog
|
||||
uses: mikepenz/release-changelog-builder-action@v5
|
||||
with:
|
||||
configurationJson: |
|
||||
{
|
||||
"template": "#{{CHANGELOG}}\n\n**Full Changelog**: #{{RELEASE_DIFF}}",
|
||||
"categories": [
|
||||
{
|
||||
"title": "## Features",
|
||||
"commits": ["^feat", "^feature"]
|
||||
},
|
||||
{
|
||||
"title": "## Bug Fixes",
|
||||
"commits": ["^fix", "^Fix"]
|
||||
},
|
||||
{
|
||||
"title": "## Other Changes",
|
||||
"commits": [".*"]
|
||||
}
|
||||
]
|
||||
}
|
||||
env:
|
||||
GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }}
|
||||
|
||||
- name: Create or update GitHub release
|
||||
env:
|
||||
GH_TOKEN: ${{ secrets.GITHUB_TOKEN }}
|
||||
TAG: ${{ github.ref_name }}
|
||||
CHANGELOG: ${{ steps.changelog.outputs.changelog }}
|
||||
run: |
|
||||
if gh release view "$TAG" >/dev/null 2>&1; then
|
||||
gh release edit "$TAG" --notes "$CHANGELOG"
|
||||
else
|
||||
gh release create "$TAG" --title "$TAG" --notes "$CHANGELOG"
|
||||
fi
|
||||
|
||||
- name: Upload Helm chart to release
|
||||
if: steps.version.outputs.is_app_release == 'true'
|
||||
env:
|
||||
GH_TOKEN: ${{ secrets.GITHUB_TOKEN }}
|
||||
TAG: ${{ github.ref_name }}
|
||||
run: gh release upload "$TAG" "$CHART_PACKAGE" --clobber
|
||||
|
||||
- name: Upload Helm chart artifact to Gitea release
|
||||
if: steps.version.outputs.is_app_release == 'true' && secrets.GITEA_TOKEN != ''
|
||||
- name: Create or update Gitea release
|
||||
if: secrets.GITEA_TOKEN != ''
|
||||
env:
|
||||
GITEA_BASE_URL: https://git.kloude.ir
|
||||
GITEA_REPO: kloude/kks-provider-plugin
|
||||
@@ -102,13 +58,9 @@ jobs:
|
||||
run: |
|
||||
set -euo pipefail
|
||||
|
||||
if [ ! -f "$CHART_PACKAGE" ]; then
|
||||
echo "Chart package not found: $CHART_PACKAGE"
|
||||
exit 1
|
||||
fi
|
||||
|
||||
api="${GITEA_BASE_URL}/api/v1/repos/${GITEA_REPO}"
|
||||
file_name=$(basename "$CHART_PACKAGE")
|
||||
release_name="$TAG"
|
||||
release_body="Release ${TAG}"
|
||||
|
||||
release_json=$(curl -fsS \
|
||||
-H "Authorization: token ${GITEA_TOKEN}" \
|
||||
@@ -119,10 +71,30 @@ jobs:
|
||||
-H "Authorization: token ${GITEA_TOKEN}" \
|
||||
-H "Content-Type: application/json" \
|
||||
"${api}/releases" \
|
||||
-d "{\"tag_name\":\"${TAG}\",\"name\":\"${TAG}\",\"target_commitish\":\"${GITHUB_SHA}\"}")
|
||||
-d "{\"tag_name\":\"${TAG}\",\"name\":\"${release_name}\",\"body\":\"${release_body}\",\"target_commitish\":\"${GITHUB_SHA}\"}")
|
||||
else
|
||||
release_id=$(python3 -c 'import json,sys; print(json.loads(sys.stdin.read())["id"])' <<< "$release_json")
|
||||
curl -fsS -X PATCH \
|
||||
-H "Authorization: token ${GITEA_TOKEN}" \
|
||||
-H "Content-Type: application/json" \
|
||||
"${api}/releases/${release_id}" \
|
||||
-d "{\"name\":\"${release_name}\",\"body\":\"${release_body}\"}" >/dev/null
|
||||
release_json=$(curl -fsS \
|
||||
-H "Authorization: token ${GITEA_TOKEN}" \
|
||||
"${api}/releases/tags/${TAG}")
|
||||
fi
|
||||
|
||||
if [ "${{ steps.version.outputs.is_app_release }}" != "true" ]; then
|
||||
exit 0
|
||||
fi
|
||||
|
||||
if [ ! -f "$CHART_PACKAGE" ]; then
|
||||
echo "Chart package not found: $CHART_PACKAGE"
|
||||
exit 1
|
||||
fi
|
||||
|
||||
release_id=$(python3 -c 'import json,sys; print(json.loads(sys.stdin.read())["id"])' <<< "$release_json")
|
||||
file_name=$(basename "$CHART_PACKAGE")
|
||||
|
||||
existing_assets=$(curl -fsS \
|
||||
-H "Authorization: token ${GITEA_TOKEN}" \
|
||||
@@ -142,24 +114,15 @@ jobs:
|
||||
--data-binary "@${CHART_PACKAGE}" \
|
||||
"${api}/releases/${release_id}/assets?name=${file_name}"
|
||||
|
||||
- name: Log in to registry
|
||||
if: steps.version.outputs.is_app_release == 'true'
|
||||
uses: docker/login-action@v3
|
||||
with:
|
||||
registry: registry.kloude.ir
|
||||
username: ${{ github.actor }}
|
||||
password: ${{ secrets.GITHUB_TOKEN }}
|
||||
|
||||
- name: Push chart to OCI registry
|
||||
if: steps.version.outputs.is_app_release == 'true'
|
||||
run: |
|
||||
owner=$(echo "${GITHUB_REPOSITORY_OWNER}" | tr '[:upper:]' '[:lower:]')
|
||||
helm push "$CHART_PACKAGE" "oci://registry.kloude.ir/kloude/kks-provider-plugin-charts"
|
||||
set -euo pipefail
|
||||
repo="oci://registry.kloude.ir/kloude-public/kks-provider-plugin-charts"
|
||||
|
||||
- name: Make OCI chart public
|
||||
if: steps.version.outputs.is_app_release == 'true'
|
||||
env:
|
||||
GH_TOKEN: ${{ secrets.GITHUB_TOKEN }}
|
||||
run: |
|
||||
owner=$(echo "${GITHUB_REPOSITORY_OWNER}" | tr '[:upper:]' '[:lower:]')
|
||||
gh api --method PATCH "/orgs/${owner}/packages/container/charts%2Fkks-provider-plugin/visibility" -f visibility=public || true
|
||||
if ! helm push --plain-http \
|
||||
--username "${{ secrets.REGISTRY_USER }}" \
|
||||
--password "${{ secrets.REGISTRY_PASS }}" \
|
||||
"$CHART_PACKAGE" "$repo"; then
|
||||
echo "OCI push failed (registry token realm/protocol issue). Continuing because chart artifact is already published to Gitea release assets."
|
||||
fi
|
||||
|
||||
@@ -19,9 +19,9 @@ Example install:
|
||||
helm install kks-provider-plugin https://github.com/kubelancloud/kks-provider-plugin/releases/download/v1.0.0/kks-provider-plugin-1.0.0.tgz \
|
||||
--namespace kube-system \
|
||||
--create-namespace \
|
||||
--set lb.serverURL=https://lb.example.kloud.team \
|
||||
--set lb.serverURL=https://lb.example.kloude.ir \
|
||||
--set lb.accessToken="$LB_TOKEN" \
|
||||
--set csi.serverURL=https://csi.example.kloud.team \
|
||||
--set csi.serverURL=https://csi.example.kloude.ir \
|
||||
--set csi.accessToken="$CSI_TOKEN"
|
||||
```
|
||||
|
||||
|
||||
@@ -1,18 +1,17 @@
|
||||
apiVersion: v2
|
||||
name: kks-provider-plugin
|
||||
description: Combined Kloud provider plugin chart (LoadBalancer + CSI) for kks clusters
|
||||
description: Combined Kloude provider plugin chart (LoadBalancer + CSI) for kks clusters
|
||||
type: application
|
||||
version: 1.1.2
|
||||
appVersion: "1.1.0"
|
||||
version: 1.2.6
|
||||
appVersion: "1.2.0"
|
||||
kubeVersion: ">=1.28.0-0"
|
||||
home: https://github.com/KubelanCloud/kks-provider-plugin
|
||||
home: https://git.kloude.ir/kloude/kks-provider-plugin
|
||||
sources:
|
||||
- https://github.com/KubelanCloud/kks-provider-plugin
|
||||
- https://git.kloude.ir/kloude/kks-provider-plugin
|
||||
keywords:
|
||||
- provider
|
||||
- loadbalancer
|
||||
- csi
|
||||
- storage
|
||||
- kloud
|
||||
- kloude
|
||||
maintainers:
|
||||
- name: Kloud Team
|
||||
- name: Kloude Tech Team
|
||||
|
||||
@@ -6,14 +6,14 @@ nameOverride: ""
|
||||
fullnameOverride: ""
|
||||
|
||||
image:
|
||||
repository: registry.kloud.team/kks/kubelancloud/kks-provider-plugin
|
||||
repository: registry.kloude.ir/kloude-public/kks-provider-plugin:latest
|
||||
tag: "latest"
|
||||
pullPolicy: IfNotPresent
|
||||
|
||||
imagePullSecrets: []
|
||||
|
||||
lb:
|
||||
serverURL: "https://lb.kloud.team"
|
||||
serverURL: "https://lb.kloude.ir"
|
||||
accessToken: ""
|
||||
existingSecret: ""
|
||||
existingSecretAccessTokenKey: access-token
|
||||
@@ -47,16 +47,16 @@ csi:
|
||||
name: storage.csi.onkksmanagement.addresslist.cloud
|
||||
sidecars:
|
||||
provisioner:
|
||||
repository: registry.kloud.team/kks/sig-storage/csi-provisioner
|
||||
repository: registry.kloude.ir/kks/sig-storage/csi-provisioner
|
||||
tag: v5.1.0
|
||||
attacher:
|
||||
repository: registry.kloud.team/kks/sig-storage/csi-attacher
|
||||
repository: registry.kloude.ir/kks/sig-storage/csi-attacher
|
||||
tag: v4.7.0
|
||||
registrar:
|
||||
repository: registry.kloud.team/kks/sig-storage/csi-node-driver-registrar
|
||||
repository: registry.kloude.ir/kks/sig-storage/csi-node-driver-registrar
|
||||
tag: v2.12.0
|
||||
livenessProbe:
|
||||
repository: registry.kloud.team/kks/sig-storage/livenessprobe
|
||||
repository: registry.kloude.ir/kks/sig-storage/livenessprobe
|
||||
tag: v2.13.1
|
||||
storageClass:
|
||||
enabled: true
|
||||
|
||||
+2
-2
@@ -93,10 +93,10 @@ func (c *Config) normalize() error {
|
||||
}
|
||||
|
||||
if c.Driver.Name == "" {
|
||||
c.Driver.Name = "storage.csi.kloud.team"
|
||||
c.Driver.Name = "storage.csi.kloude.ir"
|
||||
}
|
||||
if c.Driver.Endpoint == "" {
|
||||
c.Driver.Endpoint = "unix:///var/lib/kubelet/plugins/storage.csi.kloud.team/csi.sock"
|
||||
c.Driver.Endpoint = "unix:///var/lib/kubelet/plugins/storage.csi.kloude.ir/csi.sock"
|
||||
}
|
||||
if c.Driver.NodeID == "" {
|
||||
hostname, err := os.Hostname()
|
||||
|
||||
@@ -23,7 +23,7 @@ func TestLoadNormalizesClientDefaults(t *testing.T) {
|
||||
t.Fatalf("Load failed: %v", err)
|
||||
}
|
||||
|
||||
if cfg.Driver.Name != "storage.csi.kloud.team" {
|
||||
if cfg.Driver.Name != "storage.csi.kloude.ir" {
|
||||
t.Fatalf("unexpected driver name: %q", cfg.Driver.Name)
|
||||
}
|
||||
if cfg.Client.ServerURL != "http://192.168.84.10:9766" {
|
||||
|
||||
+2
-2
@@ -53,10 +53,10 @@ func ApplyEnvOverrides(cfg *Config) {
|
||||
cfg.Driver.Mode = CSIModeNode
|
||||
}
|
||||
if cfg.Driver.Name == "" {
|
||||
cfg.Driver.Name = "storage.csi.kloud.team"
|
||||
cfg.Driver.Name = "storage.csi.kloude.ir"
|
||||
}
|
||||
if cfg.Driver.Endpoint == "" {
|
||||
cfg.Driver.Endpoint = "unix:///var/lib/kubelet/plugins/storage.csi.kloud.team/csi.sock"
|
||||
cfg.Driver.Endpoint = "unix:///var/lib/kubelet/plugins/storage.csi.kloude.ir/csi.sock"
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -3,12 +3,12 @@
|
||||
#
|
||||
# When installed with charts/kks-provider-plugin, settings come from env vars instead.
|
||||
driver {
|
||||
name = "storage.csi.kloud.team"
|
||||
endpoint = "unix:///var/lib/kubelet/plugins/storage.csi.kloud.team/csi.sock"
|
||||
name = "storage.csi.kloude.ir"
|
||||
endpoint = "unix:///var/lib/kubelet/plugins/storage.csi.kloude.ir/csi.sock"
|
||||
mode = "all"
|
||||
}
|
||||
|
||||
client {
|
||||
server_url = "https://csi.kloud.team"
|
||||
server_url = "https://csi.kloude.ir"
|
||||
access_token = "REPLACE_WITH_CLUSTER_ACCESS_TOKEN"
|
||||
}
|
||||
|
||||
@@ -1,9 +1,11 @@
|
||||
package kloudlb
|
||||
|
||||
const (
|
||||
Finalizer = "lb.kloud.team/finalizer"
|
||||
AnnotationIP = "lb.kloud.team/ip"
|
||||
AnnotationLBID = "lb.kloud.team/id"
|
||||
Finalizer = "lb.kloude.ir/finalizer"
|
||||
AnnotationIP = "lb.kloude.ir/ip"
|
||||
AnnotationVIP = "lb.kloude.ir/vip"
|
||||
AnnotationLBID = "lb.kloude.ir/id"
|
||||
AnnotationNode = "lb.kloude.ir/node"
|
||||
|
||||
Namespace = "kube-system"
|
||||
)
|
||||
|
||||
@@ -3,7 +3,9 @@ package controller
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"net"
|
||||
"os"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
corev1 "k8s.io/api/core/v1"
|
||||
@@ -13,9 +15,9 @@ import (
|
||||
"k8s.io/client-go/informers"
|
||||
"k8s.io/client-go/kubernetes"
|
||||
corelisters "k8s.io/client-go/listers/core/v1"
|
||||
"k8s.io/client-go/tools/cache"
|
||||
"k8s.io/client-go/tools/leaderelection"
|
||||
"k8s.io/client-go/tools/leaderelection/resourcelock"
|
||||
"k8s.io/client-go/tools/cache"
|
||||
"k8s.io/client-go/util/workqueue"
|
||||
|
||||
"github.com/KubelanCloud/kks-provider-plugin/pkg/kloudlb"
|
||||
@@ -163,24 +165,39 @@ func (c *Controller) sync(ctx context.Context, key string) error {
|
||||
if ingressIP(svc) != "" {
|
||||
return nil
|
||||
}
|
||||
ports, err := servicePortsToAllocateRules(svc)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
lb, err := c.lbClient.Allocate(ctx, provisioner.AllocateRequest{
|
||||
Namespace: namespace,
|
||||
Name: name,
|
||||
Ports: ports,
|
||||
LoadBalancerSourceRanges: serviceLoadBalancerSourceRanges(svc),
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
vip := strings.TrimSpace(lb.IP)
|
||||
ingressIP, err := vipIngressIP(vip)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
patch := svc.DeepCopy()
|
||||
if patch.Annotations == nil {
|
||||
patch.Annotations = map[string]string{}
|
||||
}
|
||||
patch.Annotations[kloudlb.AnnotationIP] = lb.IP
|
||||
patch.Annotations[kloudlb.AnnotationIP] = ingressIP
|
||||
patch.Annotations[kloudlb.AnnotationVIP] = vip
|
||||
patch.Annotations[kloudlb.AnnotationLBID] = lb.ID
|
||||
if node := strings.TrimSpace(lb.Node); node != "" {
|
||||
patch.Annotations[kloudlb.AnnotationNode] = node
|
||||
}
|
||||
patch.Status = corev1.ServiceStatus{
|
||||
LoadBalancer: corev1.LoadBalancerStatus{
|
||||
Ingress: []corev1.LoadBalancerIngress{{IP: lb.IP}},
|
||||
Ingress: []corev1.LoadBalancerIngress{{IP: ingressIP}},
|
||||
},
|
||||
}
|
||||
|
||||
@@ -188,6 +205,29 @@ func (c *Controller) sync(ctx context.Context, key string) error {
|
||||
return err
|
||||
}
|
||||
|
||||
func vipIngressIP(vip string) (string, error) {
|
||||
vip = strings.TrimSpace(vip)
|
||||
if vip == "" {
|
||||
return "", fmt.Errorf("load balancer vip is empty")
|
||||
}
|
||||
if strings.Contains(vip, "/") {
|
||||
ip, _, err := net.ParseCIDR(vip)
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("invalid load balancer vip %q: %w", vip, err)
|
||||
}
|
||||
ipv4 := ip.To4()
|
||||
if ipv4 == nil {
|
||||
return "", fmt.Errorf("load balancer vip %q is not ipv4", vip)
|
||||
}
|
||||
return ipv4.String(), nil
|
||||
}
|
||||
ip := net.ParseIP(vip)
|
||||
if ip == nil || ip.To4() == nil {
|
||||
return "", fmt.Errorf("invalid load balancer vip %q", vip)
|
||||
}
|
||||
return ip.String(), nil
|
||||
}
|
||||
|
||||
func (c *Controller) finalize(ctx context.Context, svc *corev1.Service) error {
|
||||
if !containsString(svc.Finalizers, kloudlb.Finalizer) {
|
||||
return nil
|
||||
@@ -217,6 +257,47 @@ func ingressIP(svc *corev1.Service) string {
|
||||
return ""
|
||||
}
|
||||
|
||||
func servicePortsToAllocateRules(svc *corev1.Service) ([]provisioner.AllocatePortRule, error) {
|
||||
if svc == nil {
|
||||
return nil, fmt.Errorf("service is required")
|
||||
}
|
||||
if len(svc.Spec.Ports) == 0 {
|
||||
return nil, fmt.Errorf("service %s/%s has no ports", svc.Namespace, svc.Name)
|
||||
}
|
||||
out := make([]provisioner.AllocatePortRule, 0, len(svc.Spec.Ports))
|
||||
for _, p := range svc.Spec.Ports {
|
||||
protocol := string(p.Protocol)
|
||||
if protocol == "" {
|
||||
protocol = string(corev1.ProtocolTCP)
|
||||
}
|
||||
switch protocol {
|
||||
case string(corev1.ProtocolTCP), string(corev1.ProtocolUDP):
|
||||
out = append(out, provisioner.AllocatePortRule{Protocol: strings.ToLower(protocol), PortFrom: int(p.Port)})
|
||||
default:
|
||||
return nil, fmt.Errorf("service %s/%s has unsupported protocol %q for load balancer firewall", svc.Namespace, svc.Name, protocol)
|
||||
}
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
func serviceLoadBalancerSourceRanges(svc *corev1.Service) []string {
|
||||
if svc == nil || len(svc.Spec.LoadBalancerSourceRanges) == 0 {
|
||||
return nil
|
||||
}
|
||||
out := make([]string, 0, len(svc.Spec.LoadBalancerSourceRanges))
|
||||
for _, item := range svc.Spec.LoadBalancerSourceRanges {
|
||||
trimmed := strings.TrimSpace(item)
|
||||
if trimmed == "" {
|
||||
continue
|
||||
}
|
||||
out = append(out, trimmed)
|
||||
}
|
||||
if len(out) == 0 {
|
||||
return nil
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
func containsString(items []string, target string) bool {
|
||||
for _, item := range items {
|
||||
if item == target {
|
||||
|
||||
@@ -0,0 +1,64 @@
|
||||
package controller
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
corev1 "k8s.io/api/core/v1"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
)
|
||||
|
||||
func TestServicePortsToAllocateRules(t *testing.T) {
|
||||
svc := &corev1.Service{
|
||||
ObjectMeta: metav1.ObjectMeta{Namespace: "default", Name: "web"},
|
||||
Spec: corev1.ServiceSpec{Ports: []corev1.ServicePort{{Port: 80, Protocol: corev1.ProtocolTCP}, {Port: 53, Protocol: corev1.ProtocolUDP}}},
|
||||
}
|
||||
|
||||
rules, err := servicePortsToAllocateRules(svc)
|
||||
if err != nil {
|
||||
t.Fatalf("servicePortsToAllocateRules: %v", err)
|
||||
}
|
||||
if len(rules) != 2 {
|
||||
t.Fatalf("rules count = %d, want 2", len(rules))
|
||||
}
|
||||
if rules[0].Protocol != "tcp" || rules[0].PortFrom != 80 {
|
||||
t.Fatalf("unexpected first rule: %+v", rules[0])
|
||||
}
|
||||
if rules[1].Protocol != "udp" || rules[1].PortFrom != 53 {
|
||||
t.Fatalf("unexpected second rule: %+v", rules[1])
|
||||
}
|
||||
}
|
||||
|
||||
func TestServicePortsToAllocateRulesRejectsUnsupportedProtocol(t *testing.T) {
|
||||
svc := &corev1.Service{
|
||||
ObjectMeta: metav1.ObjectMeta{Namespace: "default", Name: "web"},
|
||||
Spec: corev1.ServiceSpec{Ports: []corev1.ServicePort{{Port: 80, Protocol: corev1.ProtocolSCTP}}},
|
||||
}
|
||||
|
||||
if _, err := servicePortsToAllocateRules(svc); err == nil {
|
||||
t.Fatal("expected error for unsupported protocol")
|
||||
}
|
||||
}
|
||||
|
||||
func TestServicePortsToAllocateRulesRejectsNoPorts(t *testing.T) {
|
||||
svc := &corev1.Service{ObjectMeta: metav1.ObjectMeta{Namespace: "default", Name: "web"}}
|
||||
if _, err := servicePortsToAllocateRules(svc); err == nil {
|
||||
t.Fatal("expected error for empty ports")
|
||||
}
|
||||
}
|
||||
|
||||
func TestServiceLoadBalancerSourceRanges(t *testing.T) {
|
||||
svc := &corev1.Service{
|
||||
ObjectMeta: metav1.ObjectMeta{Namespace: "default", Name: "web"},
|
||||
Spec: corev1.ServiceSpec{
|
||||
LoadBalancerSourceRanges: []string{"203.0.113.0/24", " 198.51.100.8/32 ", ""},
|
||||
},
|
||||
}
|
||||
|
||||
ranges := serviceLoadBalancerSourceRanges(svc)
|
||||
if len(ranges) != 2 {
|
||||
t.Fatalf("ranges count = %d, want 2", len(ranges))
|
||||
}
|
||||
if ranges[0] != "203.0.113.0/24" || ranges[1] != "198.51.100.8/32" {
|
||||
t.Fatalf("unexpected ranges: %#v", ranges)
|
||||
}
|
||||
}
|
||||
@@ -5,11 +5,12 @@ import (
|
||||
"fmt"
|
||||
"net"
|
||||
"os"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
corev1 "k8s.io/api/core/v1"
|
||||
coordinationv1 "k8s.io/api/coordination/v1"
|
||||
corev1 "k8s.io/api/core/v1"
|
||||
apierrors "k8s.io/apimachinery/pkg/api/errors"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"k8s.io/apimachinery/pkg/util/wait"
|
||||
@@ -141,26 +142,42 @@ func (s *Speaker) sync(ctx context.Context, key string) error {
|
||||
return err
|
||||
}
|
||||
|
||||
ip := serviceExternalIP(svc)
|
||||
if ip == "" || svc.DeletionTimestamp != nil {
|
||||
s.stopLeading(key, ip)
|
||||
vip := serviceExternalVIP(svc)
|
||||
if vip == "" || svc.DeletionTimestamp != nil {
|
||||
s.stopLeading(key, vip)
|
||||
return nil
|
||||
}
|
||||
if owner := serviceOwnerNode(svc); owner != "" && owner != s.nodeName {
|
||||
s.stopLeading(key, vip)
|
||||
return nil
|
||||
}
|
||||
|
||||
leaseName := leaseNameFor(key)
|
||||
leader, err := s.acquireLease(ctx, leaseName, ip)
|
||||
leader, err := s.acquireLease(ctx, leaseName, vip)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if leader != s.nodeName {
|
||||
s.stopLeading(key, ip)
|
||||
s.stopLeading(key, vip)
|
||||
return nil
|
||||
}
|
||||
|
||||
return s.ensureVIP(ip)
|
||||
return s.ensureVIP(vip)
|
||||
}
|
||||
|
||||
func serviceExternalIP(svc *corev1.Service) string {
|
||||
func serviceOwnerNode(svc *corev1.Service) string {
|
||||
if svc == nil || svc.Annotations == nil {
|
||||
return ""
|
||||
}
|
||||
return strings.TrimSpace(svc.Annotations[kloudlb.AnnotationNode])
|
||||
}
|
||||
|
||||
func serviceExternalVIP(svc *corev1.Service) string {
|
||||
if svc.Annotations != nil {
|
||||
if vip := strings.TrimSpace(svc.Annotations[kloudlb.AnnotationVIP]); vip != "" {
|
||||
return vip
|
||||
}
|
||||
}
|
||||
for _, ing := range svc.Status.LoadBalancer.Ingress {
|
||||
if ing.IP != "" {
|
||||
return ing.IP
|
||||
@@ -287,17 +304,15 @@ func addVIP(ifaceName, ip string) error {
|
||||
if err != nil {
|
||||
return fmt.Errorf("lookup interface %s: %w", ifaceName, err)
|
||||
}
|
||||
parsed := net.ParseIP(ip)
|
||||
if parsed == nil {
|
||||
return fmt.Errorf("invalid vip %q", ip)
|
||||
}
|
||||
addr := &netlink.Addr{
|
||||
IPNet: &net.IPNet{IP: parsed.To4(), Mask: net.CIDRMask(32, 32)},
|
||||
ipNet, err := vipToIPNet(link, ip)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
addr := &netlink.Addr{IPNet: ipNet}
|
||||
if err := netlink.AddrAdd(link, addr); err != nil {
|
||||
return fmt.Errorf("add vip %s on %s: %w", ip, ifaceName, err)
|
||||
}
|
||||
return sendGratuitousARP(link, parsed.To4())
|
||||
return sendGratuitousARP(link, ipNet.IP.To4())
|
||||
}
|
||||
|
||||
func removeVIP(ifaceName, ip string) error {
|
||||
@@ -305,16 +320,63 @@ func removeVIP(ifaceName, ip string) error {
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
parsed := net.ParseIP(ip)
|
||||
if parsed == nil {
|
||||
return fmt.Errorf("invalid vip %q", ip)
|
||||
}
|
||||
addr := &netlink.Addr{
|
||||
IPNet: &net.IPNet{IP: parsed.To4(), Mask: net.CIDRMask(32, 32)},
|
||||
ipNet, err := vipToIPNet(link, ip)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
addr := &netlink.Addr{IPNet: ipNet}
|
||||
return netlink.AddrDel(link, addr)
|
||||
}
|
||||
|
||||
func vipToIPNet(link netlink.Link, vip string) (*net.IPNet, error) {
|
||||
vip = strings.TrimSpace(vip)
|
||||
if vip == "" {
|
||||
return nil, fmt.Errorf("vip is required")
|
||||
}
|
||||
if strings.Contains(vip, "/") {
|
||||
ip, ipNet, err := net.ParseCIDR(vip)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("invalid vip %q: %w", vip, err)
|
||||
}
|
||||
ipv4 := ip.To4()
|
||||
if ipv4 == nil {
|
||||
return nil, fmt.Errorf("vip %q is not ipv4", vip)
|
||||
}
|
||||
return &net.IPNet{IP: ipv4, Mask: ipNet.Mask}, nil
|
||||
}
|
||||
|
||||
ip := net.ParseIP(vip)
|
||||
if ip == nil || ip.To4() == nil {
|
||||
return nil, fmt.Errorf("invalid vip %q", vip)
|
||||
}
|
||||
mask, err := interfaceIPv4Mask(link)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &net.IPNet{IP: ip.To4(), Mask: mask}, nil
|
||||
}
|
||||
|
||||
func interfaceIPv4Mask(link netlink.Link) (net.IPMask, error) {
|
||||
addrs, err := netlink.AddrList(link, 0)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("list ipv4 addresses on %s: %w", link.Attrs().Name, err)
|
||||
}
|
||||
for _, addr := range addrs {
|
||||
if addr.IPNet == nil || addr.IPNet.IP == nil || addr.IPNet.Mask == nil {
|
||||
continue
|
||||
}
|
||||
ip := addr.IPNet.IP.To4()
|
||||
if ip == nil {
|
||||
continue
|
||||
}
|
||||
if !ip.IsGlobalUnicast() {
|
||||
continue
|
||||
}
|
||||
return addr.IPNet.Mask, nil
|
||||
}
|
||||
return nil, fmt.Errorf("cannot infer ipv4 mask from interface %s", link.Attrs().Name)
|
||||
}
|
||||
|
||||
func sendGratuitousARP(link netlink.Link, ip net.IP) error {
|
||||
if ip == nil {
|
||||
return fmt.Errorf("ip is required")
|
||||
|
||||
@@ -5,9 +5,18 @@ type LoadBalancer struct {
|
||||
IP string `json:"ip"`
|
||||
Namespace string `json:"namespace"`
|
||||
Name string `json:"name"`
|
||||
Node string `json:"node,omitempty"`
|
||||
}
|
||||
|
||||
type AllocateRequest struct {
|
||||
Namespace string `json:"namespace"`
|
||||
Name string `json:"name"`
|
||||
Ports []AllocatePortRule `json:"ports,omitempty"`
|
||||
LoadBalancerSourceRanges []string `json:"loadBalancerSourceRanges,omitempty"`
|
||||
}
|
||||
|
||||
type AllocatePortRule struct {
|
||||
Protocol string `json:"protocol"`
|
||||
PortFrom int `json:"portFrom"`
|
||||
PortTo *int `json:"portTo,omitempty"`
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user