Add kubeflow/v1.10.0

This commit is contained in:
wbsong111
2025-06-24 12:03:10 +09:00
parent 0132496142
commit 6e8dd89e48
1531 changed files with 120984 additions and 262120 deletions
+15
View File
@@ -0,0 +1,15 @@
#!/bin/bash
set -euxo pipefail
NAMESPACES=("istio-system" "auth" "cert-manager" "oauth2-proxy" "kubeflow" "knative-serving")
for NAMESPACE in "${NAMESPACES[@]}"; do
if kubectl get namespace "$NAMESPACE" >/dev/null 2>&1; then
if [ -f "./experimental/security/PSS/static/baseline/patches/${NAMESPACE}-labels.yaml" ]; then
PATCH_OUTPUT=$(kubectl patch namespace $NAMESPACE --patch-file ./experimental/security/PSS/static/baseline/patches/${NAMESPACE}-labels.yaml 2>&1)
if echo "$PATCH_OUTPUT" | grep -q "violate the new PodSecurity"; then
exit 1
fi
fi
fi
done
+14
View File
@@ -0,0 +1,14 @@
#!/bin/bash
set -euxo pipefail
NAMESPACES=("istio-system" "auth" "cert-manager" "oauth2-proxy" "kubeflow" "knative-serving")
for NAMESPACE in "${NAMESPACES[@]}"; do
if kubectl get namespace "$NAMESPACE" >/dev/null 2>&1; then
if [ -f "./experimental/security/PSS/static/restricted/patches/${NAMESPACE}-labels.yaml" ]; then
PATCH_OUTPUT=$(kubectl patch namespace $NAMESPACE --patch-file ./experimental/security/PSS/static/restricted/patches/${NAMESPACE}-labels.yaml 2>&1)
if echo "$PATCH_OUTPUT" | grep -q "violate the new PodSecurity"; then
echo "\nWARNING PSS VIOLATED\n"
fi
fi
fi
done
@@ -0,0 +1,67 @@
#!/bin/bash
set -e
error_exit() {
echo "Error occurred in script at line: ${1}."
exit 1
}
trap 'error_exit $LINENO' ERR
echo "Install KinD..."
sudo swapoff -a
# This conditional helps running GH Workflows through
# [act](https://github.com/nektos/act)
if [ -e /swapfile ]; then
sudo rm -f /swapfile
sudo mkdir -p /tmp/etcd
sudo mount -t tmpfs tmpfs /tmp/etcd
fi
{
curl -Lo ./kind https://kind.sigs.k8s.io/dl/v0.27.0/kind-linux-amd64
chmod +x ./kind
sudo mv kind /usr/local/bin
} || { echo "Failed to install KinD"; exit 1; }
echo "Creating KinD cluster ..."
echo "
apiVersion: kind.x-k8s.io/v1alpha4
kind: Cluster
# Configure registry for KinD.
containerdConfigPatches:
- |-
[plugins.\"io.containerd.grpc.v1.cri\".registry.mirrors.\"REGISTRY_NAME:REGISTRY_PORT\"]
endpoint = [\"http://REGISTRY_NAME:REGISTRY_PORT\"]
# This is needed in order to support projected volumes with service account tokens.
# See: https://kubernetes.slack.com/archives/CEKK1KTN2/p1600268272383600
kubeadmConfigPatches:
- |
apiVersion: kubeadm.k8s.io/v1beta2
kind: ClusterConfiguration
metadata:
name: config
apiServer:
extraArgs:
\"service-account-issuer\": \"https://kubernetes.default.svc\"
\"service-account-signing-key-file\": \"/etc/kubernetes/pki/sa.key\"
nodes:
- role: control-plane
image: kindest/node:v1.32.2@sha256:f226345927d7e348497136874b6d207e0b32cc52154ad8323129352923a3142f
- role: worker
image: kindest/node:v1.32.2@sha256:f226345927d7e348497136874b6d207e0b32cc52154ad8323129352923a3142f
- role: worker
image: kindest/node:v1.32.2@sha256:f226345927d7e348497136874b6d207e0b32cc52154ad8323129352923a3142f
" | kind create cluster --config - --wait 120s
kubectl cluster-info
echo "Install Kustomize ..."
{
curl --silent --location --remote-name "https://github.com/kubernetes-sigs/kustomize/releases/download/kustomize%2Fv5.4.3/kustomize_v5.4.3_linux_amd64.tar.gz"
tar -xzvf kustomize_v5.4.3_linux_amd64.tar.gz
chmod a+x kustomize
sudo mv kustomize /usr/local/bin/kustomize
} || { echo "Failed to install Kustomize"; exit 1; }
@@ -0,0 +1,5 @@
#!/bin/bash
set -euxo pipefail
kustomize build apps/centraldashboard/upstream/overlays/kserve | kubectl apply -f -
kubectl wait --for=condition=Ready pods --all -n kubeflow --timeout=180s
@@ -3,10 +3,10 @@ set -e
echo "Installing cert-manager ..."
cd common/cert-manager
kubectl create namespace cert-manager
kustomize build cert-manager/base | kubectl apply -f -
kustomize build base | kubectl apply -f -
echo "Waiting for cert-manager to be ready ..."
kubectl wait --for=condition=ready pod -l 'app in (cert-manager,webhook)' --timeout=180s -n cert-manager
kubectl wait --for=condition=Ready pod -l 'app in (cert-manager,webhook)' --timeout=180s -n cert-manager
kubectl wait --for=jsonpath='{.subsets[0].addresses[0].targetRef.kind}'=Pod endpoints -l 'app in (cert-manager,webhook)' --timeout=180s -n cert-manager
echo "Deploy clusterissuer.cert-manager.io/kubeflow-self-signing-issuer"
+5
View File
@@ -0,0 +1,5 @@
#!/bin/bash
set -euxo pipefail
kustomize build ./common/dex/overlays/oauth2-proxy | kubectl apply -f -
kubectl wait --for=condition=Ready pods --all --timeout=180s -n auth
@@ -1,7 +1,10 @@
#!/bin/bash
set -e
echo "Installing Istio-cni ..."
cd common/istio-cni-1-22
echo "Installing Istio-cni (with ExtAuthZ from oauth2-proxy) ..."
cd common/istio-cni-1-24
kustomize build istio-crds/base | kubectl apply -f -
kustomize build istio-namespace/base | kubectl apply -f -
kustomize build istio-install/base | kubectl apply -f -
kustomize build istio-install/overlays/oauth2-proxy | kubectl apply -f -
echo "Waiting for all Istio Pods to become ready..."
kubectl wait --for=condition=Ready pods --all -n istio-system --timeout 300s
@@ -1,10 +0,0 @@
#!/bin/bash
set -e
echo "Installing Istio ..."
cd common/istio-1-22
kustomize build istio-crds/base | kubectl apply -f -
kustomize build istio-namespace/base | kubectl apply -f -
kustomize build istio-install/base | kubectl apply -f -
echo "Waiting for all Istio Pods to become ready..."
kubectl wait --for=condition=Ready pods --all -n istio-system --timeout 300s
@@ -1,17 +0,0 @@
#!/bin/bash
set -e
echo "Installing Istio configured with external authorization..."
cd common/istio-1-22
kustomize build istio-crds/base | kubectl apply -f -
kustomize build istio-namespace/base | kubectl apply -f -
kustomize build istio-install/overlays/oauth2-proxy | kubectl apply -f -
cd -
echo "Waiting for all Istio Pods to become ready..."
kubectl wait --for=condition=Ready pods --all -n istio-system --timeout=300s \
--field-selector=status.phase!=Succeeded
echo "Installing oauth2-proxy..."
cd common/oidc-client
kustomize build oauth2-proxy/overlays/m2m-self-signed/ | kubectl apply -f -
kubectl wait --for=condition=ready pod -l 'app.kubernetes.io/name=oauth2-proxy' --timeout=180s -n oauth2-proxy
+14
View File
@@ -0,0 +1,14 @@
#!/bin/bash
set -euxo pipefail
sudo apt-get update
sudo apt-get install -y apparmor-profiles
sudo apparmor_parser -R /etc/apparmor.d/usr.sbin.mysqld
cd apps/katib/upstream && kustomize build installs/katib-with-kubeflow | kubectl apply -f - && cd ../../../
kubectl wait --for=condition=Available deployment/katib-controller -n kubeflow --timeout=300s
kubectl wait --for=condition=Available deployment/katib-mysql -n kubeflow --timeout=300s
kubectl label namespace $KF_PROFILE katib.kubeflow.org/metrics-collector-injection=enabled --overwrite
@@ -1,15 +0,0 @@
#!/bin/bash
set -e
echo "Fetching KinD executable ..."
sudo swapoff -a
# This conditional helps running GH Workflows through
# [act](https://github.com/nektos/act)
if [ -e /swapfile ]; then
sudo rm -f /swapfile
sudo mkdir -p /tmp/etcd
sudo mount -t tmpfs tmpfs /tmp/etcd
fi
curl -Lo ./kind https://kind.sigs.k8s.io/dl/v0.20.0/kind-linux-amd64
chmod +x ./kind
sudo mv kind /usr/local/bin
@@ -1,14 +1,21 @@
#!/bin/bash
set -euo pipefail
echo "Installing KNative with istio-cni ..."
set +e
kustomize build common/knative/knative-serving/base | kubectl apply -f -
for ((i=1; i<=3; i++)); do
if kustomize build common/knative/knative-serving/overlays/gateways | kubectl apply -f -; then
break
fi
kubectl wait --for=condition=Ready pods --all --all-namespaces --timeout=60s --field-selector=status.phase!=Succeeded
done
set -e
kustomize build common/knative/knative-serving/base | kubectl apply -f -
kustomize build common/istio-cni-1-22/cluster-local-gateway/base | kubectl apply -f -
kustomize build common/istio-cni-1-22/kubeflow-istio-resources/base | kubectl apply -f -
kustomize build common/istio-cni-1-24/cluster-local-gateway/base | kubectl apply -f -
kustomize build common/istio-cni-1-24/kubeflow-istio-resources/base | kubectl apply -f -
kubectl wait --for=condition=Ready pods --all --all-namespaces --timeout=600s \
--field-selector=status.phase!=Succeeded
kubectl patch cm config-domain --patch '{"data":{"example.com":""}}' -n knative-serving
kubectl wait --for=condition=Ready pods --all --all-namespaces --timeout=60s --field-selector=status.phase!=Succeeded
kubectl wait --for=condition=Available deployment/activator -n knative-serving --timeout=10s
kubectl wait --for=condition=Available deployment/autoscaler -n knative-serving --timeout=10s
kubectl wait --for=condition=Available deployment/controller -n knative-serving --timeout=10s
kubectl wait --for=condition=Available deployment/webhook -n knative-serving --timeout=10s
kubectl get deployment -n knative-serving
@@ -1,14 +0,0 @@
#!/bin/bash
set -euo pipefail
echo "Installing KNative ..."
set +e
kustomize build common/knative/knative-serving/base | kubectl apply -f -
set -e
kustomize build common/knative/knative-serving/base | kubectl apply -f -
kustomize build common/istio-1-22/cluster-local-gateway/base | kubectl apply -f -
kustomize build common/istio-1-22/kubeflow-istio-resources/base | kubectl apply -f -
kubectl wait --for=condition=Ready pods --all --all-namespaces --timeout=600s \
--field-selector=status.phase!=Succeeded
kubectl patch cm config-domain --patch '{"data":{"example.com":""}}' -n knative-serving
@@ -1,15 +1,28 @@
#!/bin/bash
set -euo pipefail
set -euxo pipefail
echo "Installing Kserve ..."
cd contrib/kserve
cd apps/kserve
set +e
kustomize build kserve | kubectl apply -f -
sleep 30
kustomize build kserve | kubectl apply -f -
for ((i=1; i<=3; i++)); do
if kustomize build kserve | kubectl apply --server-side --force-conflicts -f -; then
break
fi
kubectl wait --for=condition=Ready pods --all --all-namespaces --timeout=60s --field-selector=status.phase!=Succeeded
kubectl wait --for=condition=Ready certificate/serving-cert -n kubeflow --timeout=60s
kubectl get secret kserve-webhook-server-cert -n kubeflow -o name
done
set -e
echo "Waiting for crd/clusterservingruntimes.serving.kserve.io to be available ..."
kubectl wait --for condition=established --timeout=30s crd/clusterservingruntimes.serving.kserve.io
kustomize build kserve | kubectl apply -f -
kustomize build models-web-app/overlays/kubeflow | kubectl apply -f -
kustomize build kserve | kubectl apply --server-side --force-conflicts -f -
kustomize build models-web-app/overlays/kubeflow | kubectl apply --server-side --force-conflicts -f -
kubectl wait --for=condition=Ready pods --all --all-namespaces --timeout=600s \
--field-selector=status.phase!=Succeeded
kubectl wait --for=condition=Available deployment/kserve-controller-manager -n kubeflow --timeout=10s
kubectl wait --for=condition=Available deployment/kserve-models-web-app -n kubeflow --timeout=10s
kubectl get deployment -n kubeflow -l app.kubernetes.io/name=kserve
kubectl get crd | grep -E 'inferenceservice|servingruntimes'
# Return to the original directory
cd ../../
@@ -0,0 +1,8 @@
#!/bin/bash
set -euxo
kustomize build common/user-namespace/base | kubectl apply -f -
sleep 30 # Let the profile controler reconcile the namespace
PROFILE_CONTROLLER_POD=$(kubectl get pods -n kubeflow -o json | jq -r '.items[] | select(.metadata.name | startswith("profiles-deployment")) | .metadata.name')
kubectl logs -n kubeflow "$PROFILE_CONTROLLER_POD"
KF_PROFILE=kubeflow-user-example-com
kubectl -n $KF_PROFILE get pods,configmaps,secrets
@@ -1,6 +1,7 @@
#!/bin/bash
set -e
curl --silent --location --remote-name "https://github.com/kubernetes-sigs/kustomize/releases/download/kustomize%2Fv5.2.1/kustomize_v5.2.1_linux_amd64.tar.gz"
tar -xzvf kustomize_v5.2.1_linux_amd64.tar.gz
curl --silent --location --remote-name "https://github.com/kubernetes-sigs/kustomize/releases/download/kustomize%2Fv5.4.3/kustomize_v5.4.3_linux_amd64.tar.gz"
tar -xzvf kustomize_v5.4.3_linux_amd64.tar.gz
chmod a+x kustomize
sudo mv kustomize /usr/local/bin/kustomize
@@ -7,3 +7,6 @@ kubectl -n kubeflow wait --for=condition=Ready pods -l kustomize.component=profi
echo "Installing Multitenancy Kubeflow Roles"
kustomize build common/kubeflow-roles/base | kubectl apply -f -
echo "Installing Multitenancy Network policies"
kustomize build common/networkpolicies/base | kubectl apply -f -
+13
View File
@@ -0,0 +1,13 @@
#!/bin/bash
set -e
echo "Installing oauth2-proxy..."
cd common/
kustomize build oauth2-proxy/overlays/m2m-dex-and-kind/ | kubectl apply -f -
echo "Waiting for all oauth2-proxy pods to become ready..."
kubectl wait --for=condition=Ready pod -l 'app.kubernetes.io/name=oauth2-proxy' --timeout=180s -n oauth2-proxy
echo "Waiting for all cluster-jwks-proxy pods to become ready..."
kubectl wait --for=condition=Ready pod -l 'app.kubernetes.io/name=cluster-jwks-proxy' --timeout=180s -n istio-system
kubectl wait --for=condition=Available deployment -n oauth2-proxy oauth2-proxy --timeout=180s
@@ -9,4 +9,15 @@ kustomize build env/cert-manager/platform-agnostic-multi-user | kubectl apply -f
sleep 60
kubectl wait --for=condition=Ready pods --all --all-namespaces --timeout=600s \
--field-selector=status.phase!=Succeeded
kubectl wait --for=condition=Available deployment/ml-pipeline -n kubeflow --timeout=10s
kubectl wait --for=condition=Available deployment/ml-pipeline-ui -n kubeflow --timeout=10s
kubectl wait --for=condition=Available deployment/ml-pipeline-persistenceagent -n kubeflow --timeout=10s
kubectl wait --for=condition=Available deployment/ml-pipeline-scheduledworkflow -n kubeflow --timeout=10s
kubectl wait --for=condition=Available deployment/ml-pipeline-viewer-crd -n kubeflow --timeout=10s
kubectl wait --for=condition=Available deployment/cache-server -n kubeflow --timeout=10s
kubectl wait --for=condition=Available deployment/metadata-writer -n kubeflow --timeout=10s
kubectl wait --for=condition=Available deployment/minio -n kubeflow --timeout=10s
kubectl wait --for=condition=Available deployment/mysql -n kubeflow --timeout=10s
kubectl get deployment -n kubeflow -l app=ml-pipeline
cd -
+21
View File
@@ -0,0 +1,21 @@
#!/bin/bash
set -euo pipefail
echo "Installing Pipelines ..."
kubectl apply -f apps/pipeline/upstream/third-party/metacontroller/base/crd.yaml
echo "Waiting for crd/compositecontrollers.metacontroller.k8s.io to be available ..."
kubectl wait --for condition=established --timeout=30s crd/compositecontrollers.metacontroller.k8s.io
kustomize build experimental/seaweedfs/istio | kubectl apply -f -
sleep 60
kubectl wait --for=condition=Ready pods --all --all-namespaces --timeout=600s \
--field-selector=status.phase!=Succeeded
kubectl wait --for=condition=Available deployment/ml-pipeline -n kubeflow --timeout=10s
kubectl wait --for=condition=Available deployment/ml-pipeline-ui -n kubeflow --timeout=10s
kubectl wait --for=condition=Available deployment/ml-pipeline-persistenceagent -n kubeflow --timeout=10s
kubectl wait --for=condition=Available deployment/ml-pipeline-scheduledworkflow -n kubeflow --timeout=10s
kubectl wait --for=condition=Available deployment/ml-pipeline-viewer-crd -n kubeflow --timeout=10s
kubectl wait --for=condition=Available deployment/cache-server -n kubeflow --timeout=10s
kubectl wait --for=condition=Available deployment/metadata-writer -n kubeflow --timeout=10s
kubectl wait --for=condition=Available deployment/seaweedfs -n kubeflow --timeout=10s
kubectl wait --for=condition=Available deployment/mysql -n kubeflow --timeout=10s
kubectl get deployment -n kubeflow -l app=ml-pipeline
@@ -0,0 +1,16 @@
#!/bin/bash
set -euxo
REPOSITORY_ROOT=$(git rev-parse --show-toplevel 2>/dev/null || echo "${GITHUB_WORKSPACE:-$(pwd)}")
cd "${REPOSITORY_ROOT}"
# Install Spark operator
kustomize build apps/spark/spark-operator/overlays/standalone | kubectl -n kubeflow apply --server-side -f -
# Wait for the operator controller to be ready.
kubectl -n kubeflow wait --for=condition=available --timeout=60s deploy/spark-operator-controller
kubectl -n kubeflow get pod -l app.kubernetes.io/name=spark-operator
# Wait for the operator webhook to be ready.
kubectl -n kubeflow wait --for=condition=available --timeout=30s deploy/spark-operator-webhook
kubectl -n kubeflow get pod -l app.kubernetes.io/name=spark-operator
@@ -0,0 +1,14 @@
#!/bin/bash
set -euo pipefail
cd apps/training-operator/upstream
kustomize build overlays/kubeflow | kubectl apply --server-side --force-conflicts -f -
kubectl wait --for=condition=Available deployment/training-operator -n kubeflow --timeout=180s
kubectl get deployment -n kubeflow training-operator
kubectl get pods -n kubeflow -l app=training-operator
kubectl get crd | grep -E 'tfjobs.kubeflow.org|pytorchjobs.kubeflow.org'
cd -
@@ -0,0 +1,4 @@
#!/bin/bash
set -e
curl -sfL https://raw.githubusercontent.com/aquasecurity/trivy/main/contrib/install.sh | sudo sh -s -- -b /usr/local/bin v0.48.3
@@ -0,0 +1,8 @@
#!/bin/bash
set -euxo pipefail
cd apps/volumes-web-app/upstream
kustomize build overlays/istio | kubectl apply -f -
cd ../../../
kubectl wait --for=condition=Available deployment/volumes-web-app-deployment -n kubeflow --timeout=180s
@@ -3,7 +3,7 @@
apiVersion: kubeflow.org/v1beta1
kind: Experiment
metadata:
namespace: kubeflow-user
namespace: kubeflow-user-example-com
name: grid
spec:
objective:
@@ -16,41 +16,44 @@ spec:
maxTrialCount: 2
maxFailedTrialCount: 2
parameters:
- name: lr
parameterType: double
feasibleSpace:
min: "0.01"
step: "0.005"
max: "0.05"
- name: momentum
parameterType: double
feasibleSpace:
min: "0.5"
step: "0.1"
max: "0.9"
- name: lr
parameterType: double
feasibleSpace:
min: "0.01"
step: "0.005"
max: "0.05"
- name: momentum
parameterType: double
feasibleSpace:
min: "0.5"
step: "0.1"
max: "0.9"
trialTemplate:
primaryContainerName: training-container
trialParameters:
- name: learningRate
description: Learning rate for the training model
reference: lr
- name: momentum
description: Momentum for the training model
reference: momentum
- name: learningRate
description: Learning rate for the training model
reference: lr
- name: momentum
description: Momentum for the training model
reference: momentum
trialSpec:
apiVersion: batch/v1
kind: Job
spec:
template:
metadata:
annotations:
sidecar.istio.io/inject: "false"
spec:
containers:
- name: training-container
image: docker.io/kubeflowkatib/pytorch-mnist-cpu:latest
command:
- "python3"
- "/opt/pytorch-mnist/mnist.py"
- "--epochs=1"
- "--batch-size=16"
- "--lr=${trialParameters.learningRate}"
- "--momentum=${trialParameters.momentum}"
- name: training-container
image: docker.io/kubeflowkatib/pytorch-mnist-cpu:latest
command:
- "python3"
- "/opt/pytorch-mnist/mnist.py"
- "--epochs=1"
- "--batch-size=16"
- "--lr=${trialParameters.learningRate}"
- "--momentum=${trialParameters.momentum}"
restartPolicy: Never
@@ -2,14 +2,15 @@ apiVersion: "serving.kserve.io/v1beta1"
kind: "InferenceService"
metadata:
name: "sklearn-iris"
namespace: "kubeflow-user-example-com"
spec:
predictor:
sklearn:
resources:
limits:
cpu: "1"
memory: 2Gi
requests:
cpu: "0.1"
memory: 200M
storageUri: "gs://kfserving-examples/models/sklearn/1.0/model"
limits:
cpu: "1"
memory: 2Gi
requests:
cpu: "0.1"
memory: 200M
storageUri: "gs://kfserving-examples/models/sklearn/1.0/model"
@@ -15,7 +15,7 @@ spec:
spec:
containers:
- name: test
image: kubeflownotebookswg/jupyter-scipy:v1.9.0-rc.1
image: ghcr.io/kubeflow/kubeflow/notebook-servers/jupyter-scipy:v1.10.0
imagePullPolicy: IfNotPresent
resources:
limits:
@@ -1,21 +0,0 @@
apiVersion: "kubeflow.org/v1"
kind: TFJob
metadata:
name: tfjob-simple
namespace: kubeflow
spec:
tfReplicaSpecs:
Worker:
replicas: 2
restartPolicy: OnFailure
template:
spec:
containers:
- name: tensorflow
image: gcr.io/kubeflow-ci/tf-mnist-with-summaries:1.0
command:
- "python"
- "/var/tf_mnist/mnist_with_summaries.py"
- "--log_dir=/train/logs"
- "--learning_rate=0.01"
- "--batch_size=150"
@@ -0,0 +1,40 @@
# from https://github.com/kubeflow/training-operator/blob/master/examples/pytorch/simple.yaml
# and disabled istio as stated in the documentation https://www.kubeflow.org/docs/components/training/user-guides/pytorch/
apiVersion: "kubeflow.org/v1"
kind: PyTorchJob
metadata:
name: pytorch-simple
spec:
pytorchReplicaSpecs:
Master:
replicas: 1
restartPolicy: OnFailure
template:
metadata:
labels:
sidecar.istio.io/inject: "false"
spec:
containers:
- name: pytorch
image: docker.io/kubeflowkatib/pytorch-mnist:v1beta1-45c5727
imagePullPolicy: Always
command:
- "python3"
- "/opt/pytorch-mnist/mnist.py"
- "--epochs=1"
Worker:
replicas: 1
restartPolicy: OnFailure
template:
metadata:
labels:
sidecar.istio.io/inject: "false"
spec:
containers:
- name: pytorch
image: docker.io/kubeflowkatib/pytorch-mnist:v1beta1-45c5727
imagePullPolicy: Always
command:
- "python3"
- "/opt/pytorch-mnist/mnist.py"
- "--epochs=1"
@@ -1,26 +0,0 @@
apiVersion: kind.x-k8s.io/v1alpha4
kind: Cluster
# Configure registry for KinD.
containerdConfigPatches:
- |-
[plugins."io.containerd.grpc.v1.cri".registry.mirrors."$REGISTRY_NAME:$REGISTRY_PORT"]
endpoint = ["http://$REGISTRY_NAME:$REGISTRY_PORT"]
# This is needed in order to support projected volumes with service account tokens.
# See: https://kubernetes.slack.com/archives/CEKK1KTN2/p1600268272383600
kubeadmConfigPatches:
- |
apiVersion: kubeadm.k8s.io/v1beta2
kind: ClusterConfiguration
metadata:
name: config
apiServer:
extraArgs:
"service-account-issuer": "kubernetes.default.svc"
"service-account-signing-key-file": "/etc/kubernetes/pki/sa.key"
nodes:
- role: control-plane
image: kindest/node:v1.29.4@sha256:3abb816a5b1061fb15c6e9e60856ec40d56b7b52bcea5f5f1350bc6e2320b6f8
- role: worker
image: kindest/node:v1.29.4@sha256:3abb816a5b1061fb15c6e9e60856ec40d56b7b52bcea5f5f1350bc6e2320b6f8
- role: worker
image: kindest/node:v1.29.4@sha256:3abb816a5b1061fb15c6e9e60856ec40d56b7b52bcea5f5f1350bc6e2320b6f8
@@ -0,0 +1,6 @@
{
"instances": [
[6.8, 2.8, 4.8, 1.4],
[6.0, 3.4, 4.5, 1.6]
]
}
@@ -0,0 +1,4 @@
pytest>=7.0.0
kserve>=0.15.0
kubernetes>=18.20.0
requests>=2.18.4
@@ -0,0 +1,60 @@
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
import os
from kubernetes import client
from kubernetes.client import V1ResourceRequirements
from kserve import (
constants,
KServeClient,
V1beta1InferenceService,
V1beta1InferenceServiceSpec,
V1beta1PredictorSpec,
V1beta1SKLearnSpec,
)
from utils import KSERVE_TEST_NAMESPACE
from utils import predict
def test_sklearn_kserve():
service_name = "isvc-sklearn"
predictor = V1beta1PredictorSpec(
min_replicas=1,
sklearn=V1beta1SKLearnSpec(
storage_uri="gs://kfserving-examples/models/sklearn/1.0/model",
resources=V1ResourceRequirements(
requests={"cpu": "50m", "memory": "128Mi"},
limits={"cpu": "100m", "memory": "256Mi"},
),
),
)
isvc = V1beta1InferenceService(
api_version=constants.KSERVE_V1BETA1,
kind="InferenceService",
metadata=client.V1ObjectMeta(
name=service_name, namespace=KSERVE_TEST_NAMESPACE
),
spec=V1beta1InferenceServiceSpec(predictor=predictor),
)
kserve_client = KServeClient(
config_file=os.environ.get("KUBECONFIG", "~/.kube/config")
)
kserve_client.create(isvc)
kserve_client.wait_isvc_ready(service_name, namespace=KSERVE_TEST_NAMESPACE)
res = predict(service_name, "./data/iris_input.json")
assert res["predictions"] == [1, 1]
kserve_client.delete(service_name, KSERVE_TEST_NAMESPACE)
@@ -0,0 +1,125 @@
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
import json
import logging
import os
import time
from urllib.parse import urlparse
import requests
from kubernetes import client
from kserve import KServeClient
from kserve import constants
logging.basicConfig(level=logging.INFO)
KSERVE_NAMESPACE = "kserve"
KSERVE_TEST_NAMESPACE = "kubeflow-user-example-com"
MODEL_CLASS_NAME = "modelClass"
class M2mTokenNotAvailable(Exception):
pass
def get_cluster_ip(name="istio-ingressgateway", namespace="istio-system"):
api_instance = client.CoreV1Api(client.ApiClient())
service = api_instance.read_namespaced_service(name, namespace)
if service.status.load_balancer.ingress is None:
cluster_ip = service.spec.cluster_ip
else:
if service.status.load_balancer.ingress[0].hostname:
cluster_ip = service.status.load_balancer.ingress[0].hostname
else:
cluster_ip = service.status.load_balancer.ingress[0].ip
return os.environ.get("KSERVE_INGRESS_HOST_PORT", cluster_ip)
def get_m2m_auth_token(env_name="KSERVE_M2M_TOKEN"):
try:
return os.environ[env_name]
except KeyError:
raise M2mTokenNotAvailable(env_name)
def predict(
service_name,
input_json,
protocol_version="v1",
version=constants.KSERVE_V1BETA1_VERSION,
model_name=None,
):
with open(input_json) as json_file:
data = json.load(json_file)
return predict_str(
service_name=service_name,
input_json=json.dumps(data),
protocol_version=protocol_version,
version=version,
model_name=model_name,
)
def predict_str(
service_name,
input_json,
protocol_version="v1",
version=constants.KSERVE_V1BETA1_VERSION,
model_name=None,
):
kfs_client = KServeClient(
config_file=os.environ.get("KUBECONFIG", "~/.kube/config")
)
isvc = kfs_client.get(
service_name,
namespace=KSERVE_TEST_NAMESPACE,
version=version,
)
# temporary sleep until this is fixed https://github.com/kserve/kserve/issues/604
time.sleep(10)
cluster_ip = get_cluster_ip()
host = f"{service_name}.{KSERVE_TEST_NAMESPACE}.example.com"
headers = {
"Host": host,
"Content-Type": "application/json",
}
try:
token = get_m2m_auth_token()
headers.update({"Authorization": f"Bearer {token}"})
logging.info("M2M Token Found.")
except M2mTokenNotAvailable:
logging.warn("M2M Token Not found, client authentication disabled.")
if model_name is None:
model_name = service_name
url = f"http://{cluster_ip}/v1/models/{model_name}:predict"
if protocol_version == "v2":
url = f"http://{cluster_ip}/v2/models/{model_name}/infer"
logging.info("Sending Header = %s", headers)
logging.info("Sending url = %s", url)
logging.info("Sending request data: %s", input_json)
response = requests.post(url, input_json, headers=headers)
logging.info(
"Got response code %s, content %s", response.status_code, response.content
)
if response.status_code == 200:
preds = json.loads(response.content.decode("utf-8"))
return preds
else:
response.raise_for_status()
+104
View File
@@ -0,0 +1,104 @@
#!/bin/bash
set -euxo pipefail
if kubectl get cm -n auth dex -o yaml | grep -q "redirectURIs.*http://authservice"; then
kubectl apply -f - <<EOF
apiVersion: v1
kind: ConfigMap
metadata:
name: dex
namespace: auth
data:
config.yaml: |
issuer: http://dex.auth.svc.cluster.local:5556/dex
storage:
type: kubernetes
config:
inCluster: true
web:
http: 0.0.0.0:5556
logger:
level: "debug"
format: text
oauth2:
skipApprovalScreen: true
enablePasswordDB: true
staticPasswords:
- email: user@example.com
hashFromEnv: DEX_USER_PASSWORD
username: user
userID: "15841185641784"
staticClients:
- idEnv: OIDC_CLIENT_ID
redirectURIs: ["/oauth2/callback"]
name: 'Dex Login Application'
secretEnv: OIDC_CLIENT_SECRET
EOF
kubectl rollout restart deployment -n auth dex
sleep 5
fi
# Create OIDC client secrets
kubectl create secret generic oidc-client-secret -n auth \
--from-literal=OIDC_CLIENT_ID=kubeflow-oidc-authservice \
--from-literal=OIDC_CLIENT_SECRET=pUBnBOY80SnXgjibTYM9ZWNzY2xreNGQok \
--dry-run=client -o yaml | kubectl apply -f -
kubectl create secret generic oidc-client-secret -n oauth2-proxy \
--from-literal=OIDC_CLIENT_ID=kubeflow-oidc-authservice \
--from-literal=OIDC_CLIENT_SECRET=pUBnBOY80SnXgjibTYM9ZWNzY2xreNGQok \
--dry-run=client -o yaml | kubectl apply -f -
if ! kubectl get deploy -n auth dex -o yaml | grep -q "OIDC_CLIENT_ID"; then
kubectl create secret generic dex-secret -n auth \
--from-literal=DEX_USER_PASSWORD=$(python3 -c 'from passlib.hash import bcrypt; print(bcrypt.using(rounds=12, ident="2y").hash("12341234"))') \
--dry-run=client -o yaml | kubectl apply -f -
kubectl create secret generic oidc-client-secret -n auth \
--from-literal=OIDC_CLIENT_ID=kubeflow-oidc-authservice \
--from-literal=OIDC_CLIENT_SECRET=pUBnBOY80SnXgjibTYM9ZWNzY2xreNGQok \
--dry-run=client -o yaml | kubectl apply -f -
./tests/gh-actions/install_dex.sh
fi
if kubectl get deployment -n oauth2-proxy oauth2-proxy &>/dev/null; then
kubectl rollout restart deployment -n oauth2-proxy oauth2-proxy
fi
RETRY_COUNT=0
MAX_RETRIES=3
until curl -s -o /dev/null -w "%{http_code}" http://localhost:8080/dex/health 2>/dev/null | grep -q "200\|302\|404"; do
RETRY_COUNT=$((RETRY_COUNT+1))
if [ $RETRY_COUNT -ge $MAX_RETRIES ]; then
echo "Error: Dex health endpoint not available after $MAX_RETRIES attempts"
exit 1
fi
sleep 10
done
sed -i 's/raise RuntimeError/print("ERROR:"); exit 1/g' tests/gh-actions/test_dex_login.py
# Create a temporary python script file instead of using heredoc
cat > /tmp/update_dex_login.py << 'PYTHONEOF'
import re
with open('tests/gh-actions/test_dex_login.py', 'r') as f:
content = f.read()
content = re.sub('import re', 'import re, time, sys', content, count=1)
retry_pattern = r'([ \t]+)session_cookies = dex_session_manager\.get_session_cookies\(\)'
replacement = r"""\1# Try with retries
\1for _attempt in range(3):
\1 session_cookies = dex_session_manager.get_session_cookies()
\1 if session_cookies:
\1 break
\1 if _attempt == 2: # Last attempt failed
\1 print("Error: Failed to get Dex session cookies after 3 attempts")
\1 sys.exit(1)
\1 time.sleep(5)"""
content = re.sub(retry_pattern, replacement, content, count=1)
with open('tests/gh-actions/test_dex_login.py', 'w') as f:
f.write(content)
PYTHONEOF
python3 /tmp/update_dex_login.py
rm /tmp/update_dex_login.py
+6
View File
@@ -0,0 +1,6 @@
#!/bin/bash
set -euxo pipefail
GATEWAY_SERVICE=$(kubectl get svc -n istio-system -l app=istio-ingressgateway -o jsonpath='{.items[0].metadata.name}')
nohup kubectl port-forward -n istio-system svc/$GATEWAY_SERVICE 8080:80 &
timeout 60s bash -c 'until curl -s localhost:8080 > /dev/null || curl -s -I localhost:8080 | grep -q "HTTP/"; do sleep 5; done'
@@ -93,10 +93,9 @@ for pod_name in $pod_names; do
done
done
# Exit with an error if any pod contains an error condition
if [ $error_flag -eq 1 ]; then
exit 1
fi
# This allows us to collect information about non-compliant containers without breaking the build
echo "Security check completed. Found $error_flag issues that would normally cause failure."
echo "Exiting with success for CI testing purposes."
# Exit successfully
# Always exit with success in CI environment
exit 0
+193
View File
@@ -0,0 +1,193 @@
#!/usr/bin/env python3
import re
import time
from urllib.parse import urlsplit, urlencode
import requests
import urllib3
class DexSessionManager:
"""
This is a version of the KFPClientManager() which only generates the Dex session cookies.
See https://www.kubeflow.org/docs/components/pipelines/user-guides/core-functions/connect-api/#kubeflow-platform---outside-the-cluster
"""
def __init__(
self,
endpoint_url: str,
dex_username: str,
dex_password: str,
dex_auth_type: str = "local",
skip_tls_verify: bool = False,
):
"""
Initialize the DexSessionManager
:param endpoint_url: the Kubeflow Endpoint URL
:param skip_tls_verify: if True, skip TLS verification
:param dex_username: the Dex username
:param dex_password: the Dex password
:param dex_auth_type: the auth type to use if Dex has multiple enabled, one of: ['ldap', 'local']
"""
self._endpoint_url = endpoint_url
self._skip_tls_verify = skip_tls_verify
self._dex_username = dex_username
self._dex_password = dex_password
self._dex_auth_type = dex_auth_type
self._client = None
# disable SSL verification, if requested
if self._skip_tls_verify:
urllib3.disable_warnings(urllib3.exceptions.InsecureRequestWarning)
# ensure `dex_default_auth_type` is valid
if self._dex_auth_type not in ["ldap", "local"]:
raise ValueError(
f"Invalid `dex_auth_type` '{self._dex_auth_type}', must be one of: ['ldap', 'local']"
)
def get_session_cookies(self) -> str:
"""
Get the session cookies by authenticating against Dex
:return: a string of session cookies in the form "key1=value1; key2=value2"
"""
max_retries = 3
retry_delay = 2
# Use a persistent session for cookies
session = requests.Session()
for attempt in range(max_retries):
try:
# GET the endpoint_url, which should redirect to Dex
response = session.get(
self._endpoint_url,
allow_redirects=True,
verify=not self._skip_tls_verify
)
if response.status_code == 200:
pass
elif response.status_code == 403:
# if we get 403, we might be at the oauth2-proxy sign-in page
# the default path to start the sign-in flow is `/oauth2/start?rd=<url>`
url_object = urlsplit(response.url)
url_object = url_object._replace(
path="/oauth2/start",
query=urlencode({"rd": url_object.path})
)
response = session.get(
url_object.geturl(),
allow_redirects=True,
verify=not self._skip_tls_verify
)
else:
raise RuntimeError(
f"HTTP status code '{response.status_code}' for GET against: {self._endpoint_url}"
)
# if we were NOT redirected, then the endpoint is unsecured
if len(response.history) == 0:
# No cookies are needed
return ""
# if we are at `../auth` path, we need to select an auth type
url_object = urlsplit(response.url)
if re.search(r"/auth$", url_object.path):
url_object = url_object._replace(
path=re.sub(r"/auth$", f"/auth/{self._dex_auth_type}", url_object.path)
)
# if we are at `../auth/xxxx/login` path, then we are at the login page
if re.search(r"/auth/.*/login$", url_object.path):
dex_login_url = url_object.geturl()
else:
# otherwise, we need to follow a redirect to the login page
response = session.get(
url_object.geturl(),
allow_redirects=True,
verify=not self._skip_tls_verify
)
if response.status_code != 200:
raise RuntimeError(
f"HTTP status code '{response.status_code}' for GET against: {url_object.geturl()}"
)
dex_login_url = response.url
# attempt Dex login
response = session.post(
dex_login_url,
data={"login": self._dex_username, "password": self._dex_password},
allow_redirects=True,
verify=not self._skip_tls_verify,
)
# Handle 403 specifically - might need to restart oauth flow
if response.status_code == 403:
# Try one more approach - go through the oauth2 flow again
oauth_url = f"{urlsplit(self._endpoint_url).scheme}://{urlsplit(self._endpoint_url).netloc}/oauth2/start"
response = session.get(
oauth_url,
allow_redirects=True,
verify=not self._skip_tls_verify,
)
# Continue with normal flow after restart
if response.status_code == 200 and session.cookies:
return "; ".join([f"{c.name}={c.value}" for c in session.cookies])
if response.status_code != 200:
raise RuntimeError(
f"HTTP status code '{response.status_code}' for POST against: {dex_login_url}"
)
# if we were NOT redirected, then the login credentials were probably invalid
if len(response.history) == 0:
raise RuntimeError(
f"Login credentials are probably invalid - "
f"No redirect after POST to: {dex_login_url}"
)
# if we are at `../approval` path, we need to approve the login
url_object = urlsplit(response.url)
if re.search(r"/approval$", url_object.path):
dex_approval_url = url_object.geturl()
# Approve the login
response = session.post(
dex_approval_url,
data={"approval": "approve"},
allow_redirects=True,
verify=not self._skip_tls_verify,
)
if response.status_code != 200:
raise RuntimeError(
f"HTTP status code '{response.status_code}' for POST against: {url_object.geturl()}"
)
return "; ".join([f"{c.name}={c.value}" for c in session.cookies])
except Exception as e:
if attempt == max_retries - 1: # Last attempt
print(f"All {max_retries} attempts failed. Last error: {str(e)}")
raise
print(f"Attempt {attempt + 1} failed: {str(e)}")
print(f"Retrying in {retry_delay} seconds...")
time.sleep(retry_delay)
retry_delay *= 2 # Exponential backoff
KUBEFLOW_ENDPOINT = "http://localhost:8080"
KUBEFLOW_USERNAME = "user@example.com"
KUBEFLOW_PASSWORD = "12341234"
# initialize a DexSessionManager
dex_session_manager = DexSessionManager(
endpoint_url=KUBEFLOW_ENDPOINT,
skip_tls_verify=True,
dex_username=KUBEFLOW_USERNAME,
dex_password=KUBEFLOW_PASSWORD,
dex_auth_type="local",
)
# try to get the session cookies
# NOTE: this will raise an exception if something goes wrong
session_cookies = dex_session_manager.get_session_cookies()
+70
View File
@@ -0,0 +1,70 @@
#!/bin/bash
set -euxo pipefail
NAMESPACE=${1:-kubeflow-user-example-com}
SCRIPT_DIRECTORY="$( cd "$( dirname "${BASH_SOURCE[0]}" )" && pwd )"
TEST_DIRECTORY="${SCRIPT_DIRECTORY}/kserve"
echo "=== KServe Predictor Service Labels ==="
kubectl get pods -n ${NAMESPACE} -l serving.knative.dev/service=isvc-sklearn-predictor-default --show-labels
cat <<EOF | kubectl apply -f -
apiVersion: security.istio.io/v1beta1
kind: AuthorizationPolicy
metadata:
name: allow-in-cluster-kserve
namespace: ${NAMESPACE}
spec:
rules:
- to:
- operation:
paths:
- /v1/models/*
- /v2/models/*
EOF
cat <<EOF | kubectl apply -f -
apiVersion: networking.istio.io/v1beta1
kind: VirtualService
metadata:
name: isvc-sklearn-external
namespace: ${NAMESPACE}
spec:
gateways:
- kubeflow/kubeflow-gateway
hosts:
- '*'
http:
- match:
- uri:
prefix: /kserve/${NAMESPACE}/isvc-sklearn/
rewrite:
uri: /
route:
- destination:
host: knative-local-gateway.istio-system.svc.cluster.local
headers:
request:
set:
Host: isvc-sklearn-predictor-default.${NAMESPACE}.svc.cluster.local
weight: 100
timeout: 300s
EOF
if ! command -v pytest &> /dev/null; then
echo "Installing test dependencies..."
pip install -r ${TEST_DIRECTORY}/requirements.txt
fi
export KSERVE_INGRESS_HOST_PORT=${KSERVE_INGRESS_HOST_PORT:-localhost:8080}
export KSERVE_M2M_TOKEN="$(kubectl -n ${NAMESPACE} create token default-editor)"
cd ${TEST_DIRECTORY} && pytest . -vs --log-level info
echo "=== Testing path-based external access ==="
curl -v -H "Host: isvc-sklearn.${NAMESPACE}.example.com" \
-H "Authorization: Bearer ${KSERVE_M2M_TOKEN}" \
-H "Content-Type: application/json" \
"http://${KSERVE_INGRESS_HOST_PORT}/kserve/${NAMESPACE}/isvc-sklearn/v1/models/isvc-sklearn:predict" \
-d '{"instances": [[6.8, 2.8, 4.8, 1.4], [6.0, 3.4, 4.5, 1.6]]}'
# TODO FOR FOLLOW-UP PR: Implement proper security with AuthorizationPolicy that restricts access
+49
View File
@@ -0,0 +1,49 @@
#!/usr/bin/env python3
import kfp
import sys
import time
def hello_world_op():
from kfp.components import func_to_container_op
def hello_world():
print("Hello World from Kubeflow Pipelines V1!")
return "Hello World"
return func_to_container_op(hello_world)
def hello_world_pipeline():
hello_op = hello_world_op()
hello_op()
def run_v1_pipeline(token, namespace):
client = kfp.Client(host="http://localhost:8080/pipeline", existing_token=token)
experiment = client.create_experiment("v1-pipeline-test", namespace=namespace)
pipeline_run = client.create_run_from_pipeline_func(
hello_world_pipeline,
experiment_name=experiment.name,
run_name="v1-hello-world",
namespace=namespace,
arguments={}
)
for iteration in range(15):
pipeline_status = client.get_run(pipeline_run.run_id).run.status
if pipeline_status == "Succeeded":
return
elif pipeline_status not in ["Running", "Pending"]:
sys.exit(1)
time.sleep(10)
sys.exit(1)
if __name__ == "__main__":
if len(sys.argv) != 3:
sys.exit(1)
run_v1_pipeline(sys.argv[1], sys.argv[2])
@@ -0,0 +1,96 @@
#!/usr/bin/env python3
import kfp
import sys
import time
from kfp import dsl
from kfp_server_api.exceptions import ApiException
@dsl.component
def hello_world_op() -> str:
print("Hello World from Kubeflow Pipelines V2!")
return "Hello World"
@dsl.pipeline(
name="hello-world-v2",
description="A very simple hello world pipeline"
)
def hello_world_pipeline():
hello_world_op()
def run_pipeline(token, namespace):
client = kfp.Client(host="http://localhost:8080/pipeline", existing_token=token)
try:
pipelines = client.list_pipelines()
print(f"Successfully connected to KFP server, found {len(pipelines.pipelines)} pipelines")
experiment = client.create_experiment("v2-pipeline-test", namespace=namespace)
print(f"Created experiment: v2-pipeline-test in namespace {namespace}")
run = client.create_run_from_pipeline_func(
pipeline_func=hello_world_pipeline,
experiment_name="v2-pipeline-test",
run_name="v2-test-run",
arguments={},
namespace=namespace
)
run_id = run.run_id
for _ in range(30):
status = client.get_run(run_id=run_id).state
if status == "SUCCEEDED":
return
elif status not in ["PENDING", "RUNNING"]:
print(f"Pipeline failed with status: {status}")
pods = client._get_k8s_client().list_namespaced_pod(
namespace=namespace,
label_selector=f"pipeline/runid={run_id}"
)
print(f"Found {len(pods.items)} pods for this run")
for pod in pods.items:
print(f"Pod {pod.metadata.name}: {pod.status.phase}")
sys.exit(1)
time.sleep(10)
sys.exit(1)
except Exception as exception:
print(f"Error in pipeline execution: {exception}")
sys.exit(1)
def test_unauthorized_access(token, namespace):
client = kfp.Client(host="http://localhost:8080/pipeline", existing_token=token)
try:
pipeline = client.list_runs(namespace=namespace)
sys.exit(1)
except ApiException as exception:
if exception.status != 403:
sys.exit(1)
if __name__ == "__main__":
if len(sys.argv) < 3:
sys.exit(1)
action = sys.argv[1]
token = sys.argv[2]
namespace = sys.argv[3]
if action == "run_pipeline":
run_pipeline(token, namespace)
elif action == "test_unauthorized_access":
test_unauthorized_access(token, namespace)
else:
sys.exit(1)
@@ -0,0 +1,46 @@
#!/bin/bash
set -euxo
NAMESPACE=$1
REPOSITORY_ROOT=$(git rev-parse --show-toplevel 2>/dev/null || echo "${GITHUB_WORKSPACE:-$(pwd)}")
SPARK_APPLICATION_YAML="${REPOSITORY_ROOT}/apps/spark/sparkapplication_example.yaml"
kubectl label namespace $NAMESPACE istio-injection=enabled --overwrite
kubectl get namespaces --selector=istio-injection=enabled
kubectl -n $NAMESPACE apply -f "$SPARK_APPLICATION_YAML"
# Wait for the Spark application
sleep 5
# Wait until the SparkApplication reaches the "RUNNING" state
while true; do
STATUS=$(kubectl get sparkapplication spark-pi-python -n $NAMESPACE -o jsonpath='{.status.applicationState.state}')
if [ "$STATUS" == "RUNNING" ]; then
echo "SparkApplication 'spark-pi-python' is running."
break
else
echo "Waiting for SparkApplication to be in RUNNING state. Current state: $STATUS"
sleep 5 # Check every 5 seconds
fi
done
# Wait for Spark to be ready.
sleep 5
# Wait until the Spark driver pod reaches the "Succeeded" or "Failed" phase
while true; do
POD_STATUS=$(kubectl get pod spark-pi-python-driver -n $NAMESPACE -o jsonpath='{.status.phase}')
if [ "$POD_STATUS" == "Succeeded" ] || [ "$POD_STATUS" == "Failed" ]; then
echo "Driver pod has completed with status: $POD_STATUS"
break
else
echo "Waiting for driver pod to complete. Current status: $POD_STATUS"
sleep 5 # Check every 5 seconds
fi
done
kubectl -n $NAMESPACE logs pod/spark-pi-python-driver
# Delete Spark Deployment
kubectl -n $NAMESPACE delete -f "$SPARK_APPLICATION_YAML"
+15
View File
@@ -0,0 +1,15 @@
#!/bin/bash
set -euxo pipefail
KF_PROFILE=${1:-kubeflow-user-example-com}
cat tests/gh-actions/kf-objects/training_operator_job.yaml | \
sed 's/name: pytorch-simple/name: pytorch-simple\n namespace: '"$KF_PROFILE"'/g' > /tmp/pytorch-job.yaml
kubectl apply -f /tmp/pytorch-job.yaml
kubectl wait --for=jsonpath='{.status.conditions[0].type}=Created' pytorchjob.kubeflow.org/pytorch-simple -n $KF_PROFILE --timeout=60s
kubectl get pods -n $KF_PROFILE --show-labels
kubectl wait --for=condition=Ready pod -l training.kubeflow.org/replica-type=worker -n $KF_PROFILE --timeout=180s
kubectl wait --for=condition=Succeeded pytorchjob/pytorch-simple -n $KF_PROFILE --timeout=450s
@@ -0,0 +1,98 @@
#!/bin/bash
set -euxo pipefail
KF_PROFILE=${1:-kubeflow-user-example-com}
TOKEN="$(kubectl -n $KF_PROFILE create token default-editor)"
UNAUTHORIZED_TOKEN="$(kubectl -n default create token default)"
curl --fail --show-error "localhost:8080/volumes/" -H "Authorization: Bearer ${TOKEN}" -v -c /tmp/xcrf.txt
XSRFTOKEN=$(grep XSRF-TOKEN /tmp/xcrf.txt | cut -f 7)
STORAGE_CLASS_NAME="standard"
kubectl get storageclass $STORAGE_CLASS_NAME
curl --fail --show-error \
"localhost:8080/volumes/api/storageclasses" \
-H "Authorization: Bearer ${TOKEN}" \
-H "X-XSRF-TOKEN: $XSRFTOKEN" -H "Cookie: XSRF-TOKEN=$XSRFTOKEN"
curl --fail --show-error -X POST \
"localhost:8080/volumes/api/namespaces/${KF_PROFILE}/pvcs" \
-H "Authorization: Bearer ${TOKEN}" \
-H "Content-Type: application/json" \
-H "X-XSRF-TOKEN: $XSRFTOKEN" -H "Cookie: XSRF-TOKEN=$XSRFTOKEN" \
-d "{
\"name\": \"test-pvc\",
\"namespace\": \"${KF_PROFILE}\",
\"type\": \"new\",
\"mode\": \"ReadWriteOnce\",
\"size\": \"1Gi\",
\"class\": \"${STORAGE_CLASS_NAME}\"
}"
kubectl get pvc test-pvc -n $KF_PROFILE
UNAUTHORIZED_STATUS=$(curl --silent --output /dev/null -w "%{http_code}" \
"localhost:8080/volumes/api/namespaces/${KF_PROFILE}/pvcs/test-pvc" \
-H "Authorization: Bearer ${UNAUTHORIZED_TOKEN}" \
-H "X-XSRF-TOKEN: $XSRFTOKEN" -H "Cookie: XSRF-TOKEN=$XSRFTOKEN")
if [[ "$UNAUTHORIZED_STATUS" != "403" ]]; then
echo "ERROR: Expected 403 status for unauthorized access, got $UNAUTHORIZED_STATUS"
exit 1
fi
curl --fail --show-error -X POST \
"localhost:8080/volumes/api/namespaces/${KF_PROFILE}/pvcs" \
-H "Authorization: Bearer ${TOKEN}" \
-H "Content-Type: application/json" \
-H "X-XSRF-TOKEN: $XSRFTOKEN" -H "Cookie: XSRF-TOKEN=$XSRFTOKEN" \
-d "{
\"name\": \"api-created-pvc\",
\"namespace\": \"${KF_PROFILE}\",
\"type\": \"new\",
\"mode\": \"ReadWriteOnce\",
\"size\": \"1Gi\",
\"class\": \"${STORAGE_CLASS_NAME}\"
}"
kubectl get pvc -n $KF_PROFILE
UNAUTH_DELETE_STATUS=$(curl --silent --output /dev/null -w "%{http_code}" -X DELETE \
"localhost:8080/volumes/api/namespaces/${KF_PROFILE}/pvcs/test-pvc" \
-H "Authorization: Bearer ${UNAUTHORIZED_TOKEN}" \
-H "X-XSRF-TOKEN: $XSRFTOKEN" -H "Cookie: XSRF-TOKEN=$XSRFTOKEN")
if [[ "$UNAUTH_DELETE_STATUS" != "403" ]]; then
echo "WARNING: Expected 403 status for unauthorized deletion, got $UNAUTH_DELETE_STATUS"
fi
if ! kubectl get pvc test-pvc -n $KF_PROFILE > /dev/null 2>&1; then
echo "ERROR: PVC 'test-pvc' not found after unauthorized deletion attempt"
exit 1
fi
curl --fail --show-error -X DELETE \
"localhost:8080/volumes/api/namespaces/${KF_PROFILE}/pvcs/test-pvc" \
-H "Authorization: Bearer ${TOKEN}" \
-H "X-XSRF-TOKEN: $XSRFTOKEN" -H "Cookie: XSRF-TOKEN=$XSRFTOKEN"
sleep 2
if kubectl get pvc test-pvc -n $KF_PROFILE > /dev/null 2>&1; then
echo "ERROR: PVC 'test-pvc' still exists after deletion"
exit 1
fi
curl --fail --show-error -X DELETE \
"localhost:8080/volumes/api/namespaces/${KF_PROFILE}/pvcs/api-created-pvc" \
-H "Authorization: Bearer ${TOKEN}" \
-H "X-XSRF-TOKEN: $XSRFTOKEN" -H "Cookie: XSRF-TOKEN=$XSRFTOKEN"
sleep 2
if kubectl get pvc api-created-pvc -n $KF_PROFILE > /dev/null 2>&1; then
echo "ERROR: PVC 'api-created-pvc' still exists after deletion"
exit 1
fi
echo "Test completed successfully!"
@@ -0,0 +1,401 @@
# The script:
# 1. Extract all the images used by the Kubeflow Working Groups
# - The reported image lists are saved in respective files under ../../image_lists directory
# 2. Scan the reported images using Trivy for security vulnerabilities
# - Scanned reports will be saved in JSON format inside ../../image_lists/security_scan_reports/ folder for each Working Group
# 3. The script will also generate a summary of the security scan reports with severity counts for each Working Group with images
# - Summary of security counts with images a JSON file inside ../../image_lists/summary_of_severity_counts_for_WG folder
# 4. Generate a summary of the security scan reports
# - The summary will be saved in JSON format inside ../../image_lists/summary_of_severity_counts_for_WG folder
# The script must be executed from the tests/gh-actions folder as it uses relative paths
import os
import subprocess
import re
import argparse
import json
import glob
from prettytable import PrettyTable
# Dictionary mapping Kubeflow workgroups to directories containing kustomization files
wg_dirs = {
"katib": "../../apps/katib/upstream/installs",
"pipelines": "../../apps/pipeline/upstream/env/cert-manager/platform-agnostic-multi-user",
"trainer": "../../apps/training-operator/upstream/overlays",
"manifests": "../../common/cert-manager/cert-manager/base ../../common/cert-manager/kubeflow-issuer/base ../../common/istio-1-24/istio-crds/base ../../common/istio-1-24/istio-namespace/base ../../common/istio-1-24/istio-install/overlays/oauth2-proxy ../../common/oauth2-proxy/overlays/m2m-self-signed ../../common/dex/overlays/oauth2-proxy ../../common/knative/knative-serving/overlays/gateways ../../common/knative/knative-eventing/base ../../common/istio-1-24/cluster-local-gateway/base ../../common/kubeflow-namespace/base ../../common/kubeflow-roles/base ../../common/istio-1-24/kubeflow-istio-resources/base",
"workbenches": "../../apps/pvcviewer-controller/upstream/base ../../apps/admission-webhook/upstream/overlays ../../apps/centraldashboard/overlays ../../apps/jupyter/jupyter-web-app/upstream/overlays ../../apps/volumes-web-app/upstream/overlays ../../apps/tensorboard/tensorboards-web-app/upstream/overlays ../../apps/profiles/upstream/overlays ../../apps/jupyter/notebook-controller/upstream/overlays ../../apps/tensorboard/tensorboard-controller/upstream/overlays",
"kserve": "../../apps/kserve - ../../apps/kserve/models-web-app/overlays/kubeflow",
"model-registry": "../../apps/model-registry/upstream",
"spark": "../../apps/spark/spark-operator/overlays/kubeflow",
}
DIRECTORY = "../../image_lists"
os.makedirs(DIRECTORY, exist_ok=True)
SCAN_REPORTS_DIR = os.path.join(DIRECTORY, "security_scan_reports")
ALL_SEVERITY_COUNTS = os.path.join(DIRECTORY, "severity_counts_with_images_for_WG")
SUMMARY_OF_SEVERITY_COUNTS = os.path.join(
DIRECTORY, "summary_of_severity_counts_for_WG"
)
os.makedirs(SCAN_REPORTS_DIR, exist_ok=True)
os.makedirs(ALL_SEVERITY_COUNTS, exist_ok=True)
os.makedirs(SUMMARY_OF_SEVERITY_COUNTS, exist_ok=True)
def log(*args, **kwargs):
# Custom log function that print messages with flush=True by default.
kwargs.setdefault("flush", True)
print(*args, **kwargs)
def save_images(wg, images, version):
# Saves a list of container images to a text file named after the workgroup and version.
output_file = f"../../image_lists/kf_{version}_{wg}_images.txt"
with open(output_file, "w") as f:
f.write("\n".join(images))
log(f"File {output_file} successfully created")
def validate_semantic_version(version):
# Validates a semantic version string (e.g., "0.1.2" or "latest").
regex = r"^[0-9]+\.[0-9]+\.[0-9]+$"
if re.match(regex, version) or version == "latest":
return version
else:
raise ValueError(f"Invalid semantic version: '{version}'")
def extract_images(version):
version = validate_semantic_version(version)
log(f"Running the script using Kubeflow version: {version}")
all_images = set() # Collect all unique images across workgroups
for wg, dirs in wg_dirs.items():
wg_images = set() # Collect unique images for this workgroup
for dir_path in dirs.split():
for root, _, files in os.walk(dir_path):
for file in files:
if file in [
"kustomization.yaml",
"kustomization.yml",
"Kustomization",
]:
full_path = os.path.join(root, file)
try:
# Execute `kustomize build` to render the kustomization file
result = subprocess.run(
["kustomize", "build", root],
check=True,
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
text=True,
)
except subprocess.CalledProcessError as e:
log(
f'ERROR:\t Failed "kustomize build" command for directory: {root}. See error above'
)
continue
# Use regex to find lines with 'image: <image-name>:<version>' or 'image: <image-name>'
# and '- image: <image-name>:<version>' but avoid environment variables
kustomize_images = re.findall(
r"^\s*-?\s*image:\s*([^$\s:]+(?:\:[^\s]+)?)$",
result.stdout,
re.MULTILINE,
)
wg_images.update(kustomize_images)
# Ensure uniqueness within workgroup images
uniq_wg_images = sorted(wg_images)
all_images.update(uniq_wg_images)
save_images(wg, uniq_wg_images, version)
# Ensure uniqueness across all workgroups
uniq_images = sorted(all_images)
save_images("all", uniq_images, version)
parser = argparse.ArgumentParser(
description="Extract images from Kubeflow kustomizations."
)
# Define a positional argument 'version' with optional occurrence and default value 'latest'. You can run this file as python3 <filename>.py or python <filename>.py <version>
parser.add_argument(
"version",
nargs="?",
type=str,
default="latest",
help="Kubeflow version to use (defaults to latest).",
)
args = parser.parse_args()
extract_images(args.version)
log("Started scanning images")
# Get list of text files excluding "kf_latest_all_images.txt"
files = [
f
for f in glob.glob(os.path.join(DIRECTORY, "*.txt"))
if not f.endswith("kf_latest_all_images.txt")
]
# Loop through each text file in the specified directory
for file in files:
log(f"Scanning images in {file}")
file_base_name = os.path.basename(file).replace(".txt", "")
# Directory to save reports for this specific file
file_reports_dir = os.path.join(SCAN_REPORTS_DIR, file_base_name)
os.makedirs(file_reports_dir, exist_ok=True)
# Directory to save security count
severity_count = os.path.join(file_reports_dir, "severity_counts")
os.makedirs(severity_count, exist_ok=True)
with open(file, "r") as f:
lines = f.readlines()
for line in lines:
line = line.strip()
image_name = line.split(":")[0]
image_tag = line.split(":")[1] if ":" in line else ""
image_name_scan = image_name.split("/")[-1]
if image_tag:
image_name_scan = f"{image_name_scan}_{image_tag}"
scan_output_file = os.path.join(
file_reports_dir, f"{image_name_scan}_scan.json"
)
log(f"Scanning ", line)
try:
result = subprocess.run(
[
"trivy",
"image",
"--format",
"json",
"--output",
scan_output_file,
line,
],
check=True,
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
text=True,
)
with open(scan_output_file, "r") as json_file:
scan_data = json.load(json_file)
if not scan_data.get("Results"):
log(f"No vulnerabilities found in {image_name}:{image_tag}")
else:
vulnerabilities_list = [
result["Vulnerabilities"]
for result in scan_data["Results"]
if "Vulnerabilities" in result and result["Vulnerabilities"]
]
if not vulnerabilities_list:
log(
f"The vulnerabilities detection may be insufficient because security updates are not provided for {image_name}:{image_tag}\n"
)
else:
severity_counts = {"LOW": 0, "MEDIUM": 0, "HIGH": 0, "CRITICAL": 0}
for vulnerabilities in vulnerabilities_list:
for vulnerability in vulnerabilities:
severity = vulnerability.get("Severity", "UNKNOWN")
if severity == "UNKNOWN":
continue
elif severity in severity_counts:
severity_counts[severity] += 1
report = {"image": line, "severity_counts": severity_counts}
image_table = PrettyTable()
image_table.field_names = ["Critical", "High", "Medium", "Low"]
image_table.add_row(
[
severity_counts["CRITICAL"],
severity_counts["HIGH"],
severity_counts["MEDIUM"],
severity_counts["LOW"],
]
)
log(f"{image_table}\n")
severity_report_file = os.path.join(
severity_count, f"{image_name_scan}_severity_report.json"
)
with open(severity_report_file, "w") as report_file:
json.dump(report, report_file, indent=4)
except subprocess.CalledProcessError as e:
log(f"Error scanning {image_name}:{image_tag}")
log(e.stderr)
# Combine all the JSON files into a single file with severity counts for all images
json_files = glob.glob(os.path.join(severity_count, "*.json"))
output_file = os.path.join(ALL_SEVERITY_COUNTS, f"{file_base_name}.json")
if not json_files:
log(f"No JSON files found in '{severity_count}'. Skipping combination.")
else:
combined_data = []
for json_file in json_files:
with open(json_file, "r") as jf:
combined_data.append(json.load(jf))
with open(output_file, "w") as of:
json.dump({"data": combined_data}, of, indent=4)
log(f"JSON files successfully combined into '{output_file}'")
# File to save summary of the severity counts for WGs as JSON format.
summary_file = os.path.join(
SUMMARY_OF_SEVERITY_COUNTS, "severity_summary_in_json_format.json"
)
# Initialize counters
unique_images = {} # unique set of images across all WGs
total_images = 0
total_low = 0
total_medium = 0
total_high = 0
total_critical = 0
# Initialize a dictionary to hold the final JSON data
merged_data = {}
# Loop through each JSON file in the ALL_SEVERITY_COUNTS
for file_path in glob.glob(os.path.join(ALL_SEVERITY_COUNTS, "*.json")):
# Split filename based on underscores
filename_parts = os.path.basename(file_path).split("_")
# Check if there are at least 3 parts (prefix, name, _images)
if len(filename_parts) >= 4:
# Extract name (second part)
filename = filename_parts[2]
filename = filename.capitalize()
else:
log(f"Skipping invalid filename format: {file_path}")
continue
with open(file_path, "r") as f:
data = json.load(f)["data"]
# Initialize counts for this file
image_count = len(data)
low = sum(entry["severity_counts"]["LOW"] for entry in data)
medium = sum(entry["severity_counts"]["MEDIUM"] for entry in data)
high = sum(entry["severity_counts"]["HIGH"] for entry in data)
critical = sum(entry["severity_counts"]["CRITICAL"] for entry in data)
# Update unique_images for the total counts later
for d in data:
unique_images[d["image"]] = d
# Create the output for this file
file_data = {
"images": image_count,
"LOW": low,
"MEDIUM": medium,
"HIGH": high,
"CRITICAL": critical,
}
# Update merged_data with filename as key
merged_data[filename] = file_data
# Update the total counts
unique_images = unique_images.values() # keep the set of values
total_images += len(unique_images)
total_low += sum(entry["severity_counts"]["LOW"] for entry in unique_images)
total_medium += sum(entry["severity_counts"]["MEDIUM"] for entry in unique_images)
total_high += sum(entry["severity_counts"]["HIGH"] for entry in unique_images)
total_critical += sum(entry["severity_counts"]["CRITICAL"] for entry in unique_images)
# Add total counts to merged_data
merged_data["total"] = {
"images": total_images,
"LOW": total_low,
"MEDIUM": total_medium,
"HIGH": total_high,
"CRITICAL": total_critical,
}
log("Summary in Json Format:")
log(json.dumps(merged_data, indent=4))
# Write the final output to a file
with open(summary_file, "w") as summary_f:
json.dump(merged_data, summary_f, indent=4)
log(f"Summary written to: {summary_file} as JSON format")
# Load JSON content from the file
with open(summary_file, "r") as file:
data = json.load(file)
# Define a mapping for working group names
working_group_name_mapping = {
"Katib": "Katib",
"Pipelines": "Pipelines",
"Workbenches": "Workbenches(Notebooks)",
"Kserve": "Kserve",
"Manifests": "Manifests",
"Trainer": "Trainer",
"Model-registry": "Model Registry",
"Spark": "Spark",
"total": "All Images",
}
# Create PrettyTable
summary_table = PrettyTable()
summary_table.field_names = [
"Working Group",
"Images",
"Critical CVE",
"High CVE",
"Medium CVE",
"Low CVE",
]
# Populate the table with data
for working_group_key in working_group_name_mapping:
if working_group_key in data: # Check if the working group exists in the data
working_group_data = data[working_group_key]
summary_table.add_row(
[
working_group_name_mapping[working_group_key],
working_group_data["images"],
working_group_data["CRITICAL"],
working_group_data["HIGH"],
working_group_data["MEDIUM"],
working_group_data["LOW"],
]
)
# log the table
log(summary_table)
# Write the table output to a file in the specified folder
summary_table_output_file = (
SUMMARY_OF_SEVERITY_COUNTS + "/summary_of_severity_counts_for_WGs_in_table.txt"
)
with open(summary_table_output_file, "w") as file:
file.write(str(summary_table))
log("Output saved to:", summary_table_output_file)
log("Severity counts with images respect to WGs are saved in the", ALL_SEVERITY_COUNTS)
log("Scanned JSON reports on images are saved in", SCAN_REPORTS_DIR)
@@ -1,41 +0,0 @@
#!/bin/bash
CRONJOB_NAME=kubeflow-m2m-oidc-configurator
NAMESPACE=istio-system
# Function to get the latest Job created by the CronJob
get_latest_job() {
kubectl get jobs -n "${NAMESPACE}" \
--sort-by=.metadata.creationTimestamp -o json \
| jq --arg cronjob_name "${CRONJOB_NAME}" -r '.items[] | select(.metadata.ownerReferences[] | select(.name==$cronjob_name)) | .metadata.name' \
| tail -n 1
}
# Wait until a Job is created
echo "Waiting for a Job to be created by the ${CRONJOB_NAME} CronJob..."
while true; do
JOB_NAME=$(get_latest_job)
if [[ -n "${JOB_NAME}" ]]; then
echo "Job ${JOB_NAME} created."
break
fi
sleep 5
echo "Waiting..."
done
# Wait for the Job to complete successfully
echo "Waiting for the Job ${JOB_NAME} to complete..."
while true; do
STATUS=$(kubectl get job "${JOB_NAME}" -n "${NAMESPACE}" -o jsonpath='{.status.conditions[?(@.type=="Complete")].status}')
if [[ "${STATUS}" == "True" ]]; then
echo "Job ${JOB_NAME} completed successfully."
break
fi
FAILED=$(kubectl get job "${JOB_NAME}" -n "${NAMESPACE}" -o jsonpath='{.status.conditions[?(@.type=="Failed")].status}')
if [[ "${FAILED}" == "True" ]]; then
echo "Job ${JOB_NAME} failed."
exit 1
fi
sleep 5
done