Skip to content
Merged
Show file tree
Hide file tree
Changes from all 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"
97 changes: 77 additions & 20 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,53 +90,58 @@ def cluster_exists(self, cluster_name: str) -> bool:
return False
raise

def _resolve_secret_ref(self, ref: str) -> str:
"""Verify that *ref* is a Vault-synced K8s Secret and return its name.
@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) -> client.V1Secret:
"""Verify that *ref* is a Vault-synced K8s Secret and return it.

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

def resolve_secret_refs(self, refs: list[str]) -> list[str]:
"""Resolve a list of secretRefs to their K8s Secret names."""
return [self._resolve_secret_ref(r) for r in refs]
return [self._resolve_secret_ref(r).metadata.name for r in refs]

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)"
)
source = self._resolve_secret_ref(ref)
secret_name = source.metadata.name

keys = sorted((source.data or {}).keys())
copied_name = f"{fjob_name}-{ref}"
Expand Down Expand Up @@ -116,7 +172,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