Skip to content
Merged
Show file tree
Hide file tree
Changes from 3 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions config/forge/workflows/tasks.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,8 @@ spec:
value: "$(params.env)"
- name: FOURNOS_STEP
value: "$(params.job-step)"
- name: FOURNOS_SECRETS
value: /var/run/secrets/fournos
volumeMounts:
- name: kubeconfig
mountPath: /var/run/secrets/fournos-kubeconfig
Expand Down
2 changes: 1 addition & 1 deletion dev/mock-resolve/resolve.sh
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,6 @@ echo "[mock-resolve] setting secretRefs"
kubectl patch fournosjob "${FOURNOS_JOB_NAME}" \
-n "${FOURNOS_NAMESPACE}" \
--type=merge \
-p '{"spec":{"secretRefs":["vault-placeholder"]}}'
-p '{"spec":{"secretRefs":["placeholder"]}}'

echo "[mock-resolve] done"
4 changes: 3 additions & 1 deletion dev/mock-secrets.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -40,13 +40,15 @@ stringData:

---
# Mock Vault-synced secret for the secret-demo sample job.
# Uses managed-by: fournos-mock (not fournos-vault-sync) so the
# destructive sync in sync_vault_secrets.py won't garbage-collect it.
apiVersion: v1
kind: Secret
metadata:
name: vault-placeholder
labels:
fournos.dev/vault-entry: "true"
app.kubernetes.io/managed-by: fournos-vault-sync
app.kubernetes.io/managed-by: fournos-mock
type: Opaque
stringData:
placeholder: "mock-secret-value"
93 changes: 76 additions & 17 deletions fournos/core/clusters.py
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,57 @@ def resolve_kubeconfig_secret(self, cluster_name: str) -> str:
"""Return the Secret name that holds the kubeconfig for *cluster_name*."""
return settings.kubeconfig_secret_pattern.format(cluster=cluster_name)

def copy_kubeconfig_secret(
self, cluster_name: str, fjob_name: str, owner_ref: dict
) -> str:
"""Copy the kubeconfig Secret for *cluster_name* into the operator namespace.

Returns the name of the copied Secret (``<fjob_name>-kubeconfig``).
Idempotent: a 409 (AlreadyExists) is silently ignored.
"""
source_name = self.resolve_kubeconfig_secret(cluster_name)
source = self._k8s.read_namespaced_secret(
source_name, settings.secrets_namespace
)

copied_name = f"{fjob_name}-kubeconfig"

copy_body = client.V1Secret(
metadata=client.V1ObjectMeta(
name=copied_name,
namespace=settings.namespace,
labels={LABEL_MANAGED_BY: "fournos"},
owner_references=[
client.V1OwnerReference(
api_version=owner_ref["apiVersion"],
kind=owner_ref["kind"],
name=owner_ref["name"],
uid=owner_ref["uid"],
controller=False,
block_owner_deletion=True,
)
],
),
type=source.type,
data=source.data,
)

try:
self._k8s.create_namespaced_secret(settings.namespace, copy_body)
logger.info(
"Copied kubeconfig %s from %s as %s",
source_name,
settings.secrets_namespace,
copied_name,
)
except client.exceptions.ApiException as exc:
if exc.status == 409:
logger.debug("Kubeconfig copy %s already exists (409)", copied_name)
else:
raise

return copied_name

def cluster_exists(self, cluster_name: str) -> bool:
"""Return True if the kubeconfig Secret for *cluster_name* exists."""
secret_name = self.resolve_kubeconfig_secret(cluster_name)
Expand All @@ -39,33 +90,43 @@ def cluster_exists(self, cluster_name: str) -> bool:
return False
raise

@staticmethod
def _vault_secret_name(ref: str) -> str:
"""Apply the vault secret naming pattern to a user-supplied ref."""
return settings.vault_secret_pattern.format(entry=ref)

def _resolve_secret_ref(self, ref: str) -> str:
"""Verify that *ref* is a Vault-synced K8s Secret and return its name.

The Secret is read from ``secrets_namespace``. The
``fournos.dev/vault-entry=true`` label is checked to confirm
Users supply refs without the ``vault-`` prefix; the pattern from
``settings.vault_secret_pattern`` is applied to derive the real
Secret name. The Secret is read from ``secrets_namespace``.
The ``fournos.dev/vault-entry=true`` label is checked to confirm
the Secret was actually imported from Vault.

Raises ``KeyError`` if the Secret does not exist or is not
a Vault-synced secret.
"""
secret_name = self._vault_secret_name(ref)
try:
secret = self._k8s.read_namespaced_secret(ref, settings.secrets_namespace)
secret = self._k8s.read_namespaced_secret(
secret_name, settings.secrets_namespace
)
except client.exceptions.ApiException as exc:
if exc.status == 404:
raise KeyError(
f"Secret {ref!r} not found in namespace "
f"{settings.secrets_namespace}"
f"Secret {secret_name!r} (ref {ref!r}) not found in "
f"namespace {settings.secrets_namespace}"
) from exc
raise
labels = secret.metadata.labels or {}
if labels.get(LABEL_VAULT_ENTRY) != "true":
raise KeyError(
f"Secret {ref!r} exists but is not a Vault-synced secret "
f"Secret {secret_name!r} exists but is not a Vault-synced secret "
f"(missing {LABEL_VAULT_ENTRY}=true label)"
)
logger.debug("Validated secretRef %s", ref)
return ref
logger.debug("Validated secretRef %s -> %s", ref, secret_name)
return secret_name

def resolve_secret_refs(self, refs: list[str]) -> list[str]:
"""Resolve a list of secretRefs to their K8s Secret names."""
Expand All @@ -74,18 +135,15 @@ def resolve_secret_refs(self, refs: list[str]) -> list[str]:
def copy_secret(self, ref: str, fjob_name: str, owner_ref: dict) -> ResolvedSecret:
"""Copy a Vault-synced Secret from the secrets namespace into the pod namespace.

*ref* is the user-supplied name (without ``vault-`` prefix).
The copy is named ``<fjob_name>-<ref>`` and carries an ownerReference
back to the FournosJob so K8s GC cleans it up automatically.
Idempotent: a 409 (AlreadyExists) is silently ignored.
"""
source = self._k8s.read_namespaced_secret(ref, settings.secrets_namespace)

labels = source.metadata.labels or {}
if labels.get(LABEL_VAULT_ENTRY) != "true":
raise KeyError(
f"Secret {ref!r} in {settings.secrets_namespace} is not a "
f"Vault-synced secret (missing {LABEL_VAULT_ENTRY}=true label)"
)
secret_name = self._resolve_secret_ref(ref)
source = self._k8s.read_namespaced_secret(
secret_name, settings.secrets_namespace
)

keys = sorted((source.data or {}).keys())
copied_name = f"{fjob_name}-{ref}"
Expand Down Expand Up @@ -116,7 +174,8 @@ def copy_secret(self, ref: str, fjob_name: str, owner_ref: dict) -> ResolvedSecr
try:
self._k8s.create_namespaced_secret(settings.namespace, copy_body)
logger.info(
"Copied secret %s from %s as %s",
"Copied secret %s (ref %s) from %s as %s",
secret_name,
ref,
settings.secrets_namespace,
copied_name,
Expand Down
22 changes: 20 additions & 2 deletions fournos/handlers/execution.py
Original file line number Diff line number Diff line change
Expand Up @@ -147,7 +147,25 @@ def reconcile_admitted(spec, name, namespace, status, patch, body):

if pr is None:
cluster = status.get("cluster", "")
secret = ctx.registry.resolve_kubeconfig_secret(cluster)

try:
kubeconfig_secret = ctx.registry.copy_kubeconfig_secret(
cluster, name, owner_ref(body)
)
except client.exceptions.ApiException as exc:
patch.status["phase"] = Phase.FAILED
patch.status["message"] = f"Failed to copy kubeconfig: {exc.reason}"
set_condition(
patch,
conditions,
COND_PIPELINE_RUN_READY,
"False",
"KubeconfigNotFound",
f"Failed to copy kubeconfig: {exc.reason}",
)
ctx.kueue.delete_workload(name)
logger.error("Job %s: kubeconfig copy failed: %s", name, exc)
return

hardware = spec.get("hardware") or {}
gpu_count = hardware.get("gpuCount", 0)
Expand Down Expand Up @@ -183,7 +201,7 @@ def reconcile_admitted(spec, name, namespace, status, patch, body):
forge_project=spec["forge"]["project"],
forge_config=spec["forge"],
env=spec.get("env", {}),
kubeconfig_secret=secret,
kubeconfig_secret=kubeconfig_secret,
gpu_count=gpu_count,
resolved_secrets=resolved_secrets,
cluster=cluster,
Expand Down
11 changes: 6 additions & 5 deletions manifests/crd.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -101,11 +101,12 @@ spec:
secretRefs:
type: array
description: >-
Vault-synced K8s Secret names (vault-<entry>) to mount
into the pipeline. Each name must correspond to a
Kubernetes Secret with the fournos.dev/vault-entry=true
label. Populated by Forge during the Resolving phase
when not provided by the user.
Vault entry names (without the vault- prefix) to mount
into the pipeline. The operator prepends vault- to
look up the corresponding K8s Secret (which must carry
the fournos.dev/vault-entry=true label) in the secrets
namespace. Populated by Forge during the Resolving
phase when not provided by the user.
items:
type: string
pattern: "^[a-z0-9]([a-z0-9\\-]{0,61}[a-z0-9])?$"
Expand Down
2 changes: 1 addition & 1 deletion tests/conftest.py
Original file line number Diff line number Diff line change
Expand Up @@ -464,7 +464,7 @@ def create_stale_pipelinerun(k8s, name: str) -> None:
"""),
},
{"name": "env", "value": ""},
{"name": "kubeconfig-secret", "value": "kubeconfig-cluster-1"},
{"name": "kubeconfig-secret", "value": f"{name}-kubeconfig"},
{"name": "gpu-count", "value": "0"},
],
},
Expand Down
27 changes: 23 additions & 4 deletions tests/test_scheduling.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@

from fournos.core.constants import Phase
from tests.conftest import (
NAMESPACE,
create_job,
get_job,
get_k8s_resource,
Expand Down Expand Up @@ -43,8 +44,26 @@ def test_cluster_pinned(k8s):
flavor = get_workload_flavor("test-cluster")
assert flavor == "cluster-2", f"Workload flavor should be cluster-2, got {flavor!r}"
secret = get_pipelinerun_param("test-cluster", "kubeconfig-secret")
assert secret == "kubeconfig-cluster-2", (
f"PipelineRun kubeconfig-secret should be kubeconfig-cluster-2, got {secret!r}"
assert secret == "test-cluster-kubeconfig", (
f"PipelineRun kubeconfig-secret should be test-cluster-kubeconfig, got {secret!r}"
)

kc = get_k8s_resource("secret", "test-cluster-kubeconfig")
assert "kubeconfig" in (kc.get("data") or {}), (
f"Copied kubeconfig secret should have a 'kubeconfig' key, got {list((kc.get('data') or {}).keys())}"
)
kc_owners = kc.get("metadata", {}).get("ownerReferences", [])
assert any(
o.get("kind") == "FournosJob" and o.get("name") == "test-cluster"
for o in kc_owners
), f"Copied kubeconfig should have FournosJob ownerRef, got {kc_owners!r}"
assert kc.get("metadata", {}).get("namespace") == NAMESPACE, (
"Copied kubeconfig should be in the operator namespace"
)

refs = get_pipelinerun_param("test-cluster", "secret-refs")
assert refs == ["placeholder"], (
f"Mock resolver should set secret-refs to ['placeholder'], got {refs!r}"
)

phase = poll_phase(
Expand Down Expand Up @@ -115,8 +134,8 @@ def test_cluster_and_hardware(k8s):
flavor = get_workload_flavor("test-cluster-hw")
assert flavor == "cluster-4", f"Workload flavor should be cluster-4, got {flavor!r}"
secret = get_pipelinerun_param("test-cluster-hw", "kubeconfig-secret")
assert secret == "kubeconfig-cluster-4", (
f"PipelineRun kubeconfig-secret should be kubeconfig-cluster-4, got {secret!r}"
assert secret == "test-cluster-hw-kubeconfig", (
f"PipelineRun kubeconfig-secret should be test-cluster-hw-kubeconfig, got {secret!r}"
)

phase = poll_phase(
Expand Down
Loading
Loading