diff --git a/.github/workflows/build-operator.yaml b/.github/workflows/build-operator.yaml index 08f9e57d..6f4d0288 100644 --- a/.github/workflows/build-operator.yaml +++ b/.github/workflows/build-operator.yaml @@ -1,4 +1,4 @@ -name: Build & Push Streaming Operator Image +name: Publish Streaming Operator on: push: @@ -37,3 +37,43 @@ jobs: google_ar_image_name: ${{ matrix.image }} google_workload_identity_provider: projects/868781662168/locations/global/workloadIdentityPools/prod-github/providers/github-oidc-pool google_service_account: gha-gcr-push@sac-prod-sa.iam.gserviceaccount.com + + publish-chart: + name: Publish Helm chart + needs: build-and-push + runs-on: ubuntu-latest + concurrency: + group: github-pages + cancel-in-progress: false + permissions: + contents: write + env: + CHART_DIR: sentry_streams_k8s/chart/streaming-operator + CHART_VERSION: 0.0.0-${{ github.sha }} + IMAGE_TAG: ${{ github.sha }} + steps: + - uses: actions/checkout@11bd71901bbe5b1630ceea73d27597364c9af683 + with: + fetch-depth: 0 + + - name: Prepare and lint chart + run: | + sed -i "s/^version:.*/version: ${CHART_VERSION}/" "${CHART_DIR}/Chart.yaml" + sed -i "s/^appVersion:.*/appVersion: \"${IMAGE_TAG}\"/" "${CHART_DIR}/Chart.yaml" + helm lint "${CHART_DIR}" --set workloadNamespace=streaming-pipelines + helm template streaming-operator "${CHART_DIR}" \ + --set workloadNamespace=streaming-pipelines > /dev/null + + - name: Configure Git + run: | + git config user.name "$GITHUB_ACTOR" + git config user.email "$GITHUB_ACTOR@users.noreply.github.com" + + - name: Release chart + uses: helm/chart-releaser-action@cae68fefc6b5f367a0275617c9f83181ba54714f # v1.7.0 + with: + charts_dir: sentry_streams_k8s/chart + skip_existing: true + mark_as_latest: false + env: + CR_TOKEN: ${{ secrets.GITHUB_TOKEN }} diff --git a/.github/workflows/ci.yaml b/.github/workflows/ci.yaml index fd61d2f9..fedc9f6f 100644 --- a/.github/workflows/ci.yaml +++ b/.github/workflows/ci.yaml @@ -30,6 +30,18 @@ jobs: - name: Run k8s test run: make tests-k8s + chart: + name: "Lint streaming-operator helm chart" + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@f43a0e5ff2bd294095638e18286ca9a3d1956744 # v3 + - name: Lint and render chart + run: | + helm lint sentry_streams_k8s/chart/streaming-operator \ + --set workloadNamespace=streaming-pipelines + helm template streaming-operator sentry_streams_k8s/chart/streaming-operator \ + --set workloadNamespace=streaming-pipelines > /dev/null + integration-tests: name: "Run integration tests" runs-on: ubuntu-latest diff --git a/.github/workflows/docs.yaml b/.github/workflows/docs.yaml index f7c417dd..bb0e3abb 100644 --- a/.github/workflows/docs.yaml +++ b/.github/workflows/docs.yaml @@ -9,6 +9,9 @@ jobs: docs: name: Sphinx runs-on: ubuntu-latest + concurrency: + group: github-pages + cancel-in-progress: false steps: - uses: actions/checkout@ee0669bd1cc54295c223e0bb666b733df41de1c5 # v2 - uses: actions/setup-python@a26af69be951a213d495a4c3e4e4022e16d87065 # v5 @@ -23,7 +26,7 @@ jobs: with: github_token: ${{ secrets.GITHUB_TOKEN }} publish_dir: sentry_streams/docs/build - force_orphan: true + keep_files: true - name: Archive Docs uses: actions/upload-artifact@ea165f8d65b6e75b540449e92b4886f43607fa02 # v4 diff --git a/sentry_streams_k8s/chart/streaming-operator/Chart.yaml b/sentry_streams_k8s/chart/streaming-operator/Chart.yaml new file mode 100644 index 00000000..cef5dbe8 --- /dev/null +++ b/sentry_streams_k8s/chart/streaming-operator/Chart.yaml @@ -0,0 +1,8 @@ +apiVersion: v2 +name: streaming-operator +type: application +kubeVersion: ">=1.23.0-0" +version: 0.0.0 +appVersion: "latest" +sources: + - https://github.com/getsentry/streams/tree/main/sentry_streams_k8s diff --git a/sentry_streams_k8s/sentry_streams_k8s/operator/manifests/crd.yaml b/sentry_streams_k8s/chart/streaming-operator/crds/crd.yaml similarity index 100% rename from sentry_streams_k8s/sentry_streams_k8s/operator/manifests/crd.yaml rename to sentry_streams_k8s/chart/streaming-operator/crds/crd.yaml diff --git a/sentry_streams_k8s/chart/streaming-operator/templates/_helpers.tpl b/sentry_streams_k8s/chart/streaming-operator/templates/_helpers.tpl new file mode 100644 index 00000000..22e4be70 --- /dev/null +++ b/sentry_streams_k8s/chart/streaming-operator/templates/_helpers.tpl @@ -0,0 +1,30 @@ +{{- define "streaming-operator.name" -}} +{{- .Chart.Name | trunc 63 | trimSuffix "-" -}} +{{- end -}} + +{{- define "streaming-operator.fullname" -}} +{{- include "streaming-operator.name" . -}} +{{- end -}} + +{{- define "streaming-operator.selectorLabels" -}} +app.kubernetes.io/name: {{ include "streaming-operator.name" . }} +app.kubernetes.io/instance: {{ .Release.Name }} +{{- end -}} + +{{- define "streaming-operator.labels" -}} +{{ include "streaming-operator.selectorLabels" . }} +app.kubernetes.io/version: {{ .Chart.AppVersion | quote }} +app.kubernetes.io/managed-by: {{ .Release.Service }} +helm.sh/chart: {{ printf "%s-%s" .Chart.Name .Chart.Version }} +{{- with .Values.labels }} +{{ toYaml . }} +{{- end }} +{{- end -}} + +{{- define "streaming-operator.serviceAccountName" -}} +{{- default (include "streaming-operator.fullname" .) .Values.serviceAccount.name -}} +{{- end -}} + +{{- define "streaming-operator.workloadNamespace" -}} +{{- required "workloadNamespace must be set" .Values.workloadNamespace -}} +{{- end -}} diff --git a/sentry_streams_k8s/chart/streaming-operator/templates/deployment.yaml b/sentry_streams_k8s/chart/streaming-operator/templates/deployment.yaml new file mode 100644 index 00000000..a419c550 --- /dev/null +++ b/sentry_streams_k8s/chart/streaming-operator/templates/deployment.yaml @@ -0,0 +1,64 @@ +apiVersion: apps/v1 +kind: Deployment +metadata: + name: {{ include "streaming-operator.fullname" . }} + namespace: {{ .Release.Namespace }} + labels: + {{- include "streaming-operator.labels" . | nindent 4 }} +spec: + replicas: 1 + strategy: + type: Recreate + selector: + matchLabels: + {{- include "streaming-operator.selectorLabels" . | nindent 6 }} + template: + metadata: + labels: + {{- include "streaming-operator.labels" . | nindent 8 }} + {{- with .Values.podLabels }} + {{- toYaml . | nindent 8 }} + {{- end }} + {{- with .Values.podAnnotations }} + annotations: + {{- toYaml . | nindent 8 }} + {{- end }} + spec: + serviceAccountName: {{ include "streaming-operator.serviceAccountName" . }} + securityContext: + runAsNonRoot: true + runAsUser: 1000 + runAsGroup: 1000 + seccompProfile: + type: RuntimeDefault + containers: + - name: operator + image: "{{ .Values.image.repository }}:{{ .Values.image.tag | default .Chart.AppVersion }}" + imagePullPolicy: {{ .Values.image.pullPolicy }} + env: + - name: WORKLOAD_NAMESPACE + value: {{ include "streaming-operator.workloadNamespace" . | quote }} + {{- with .Values.env }} + {{- toYaml . | nindent 12 }} + {{- end }} + resources: + {{- toYaml .Values.resources | nindent 12 }} + securityContext: + allowPrivilegeEscalation: false + capabilities: + drop: [ALL] + {{- with .Values.nodeSelector }} + nodeSelector: + {{- toYaml . | nindent 8 }} + {{- end }} + {{- with .Values.tolerations }} + tolerations: + {{- toYaml . | nindent 8 }} + {{- end }} + {{- with .Values.affinity }} + affinity: + {{- toYaml . | nindent 8 }} + {{- end }} + {{- with .Values.priorityClassName }} + priorityClassName: {{ . }} + {{- end }} diff --git a/sentry_streams_k8s/chart/streaming-operator/templates/namespace.yaml b/sentry_streams_k8s/chart/streaming-operator/templates/namespace.yaml new file mode 100644 index 00000000..2fddbe27 --- /dev/null +++ b/sentry_streams_k8s/chart/streaming-operator/templates/namespace.yaml @@ -0,0 +1,8 @@ +apiVersion: v1 +kind: Namespace +metadata: + name: {{ include "streaming-operator.workloadNamespace" . }} + annotations: + helm.sh/resource-policy: keep + labels: + {{- include "streaming-operator.labels" . | nindent 4 }} diff --git a/sentry_streams_k8s/chart/streaming-operator/templates/rbac.yaml b/sentry_streams_k8s/chart/streaming-operator/templates/rbac.yaml new file mode 100644 index 00000000..dfa6be89 --- /dev/null +++ b/sentry_streams_k8s/chart/streaming-operator/templates/rbac.yaml @@ -0,0 +1,68 @@ +apiVersion: rbac.authorization.k8s.io/v1 +kind: ClusterRole +metadata: + name: {{ include "streaming-operator.fullname" . }} + labels: + {{- include "streaming-operator.labels" . | nindent 4 }} +rules: + - apiGroups: [apiextensions.k8s.io] + resources: [customresourcedefinitions] + verbs: [list, watch] + - apiGroups: [""] + resources: [namespaces] + verbs: [list, watch] + - apiGroups: [streams.sentry.io] + resources: [streamingpipelines] + verbs: [list, watch, patch] + - apiGroups: [streams.sentry.io] + resources: [streamingpipelines/status] + verbs: [patch] + - apiGroups: [""] + resources: [events] + verbs: [create] +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: ClusterRoleBinding +metadata: + name: {{ include "streaming-operator.fullname" . }} + labels: + {{- include "streaming-operator.labels" . | nindent 4 }} +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: ClusterRole + name: {{ include "streaming-operator.fullname" . }} +subjects: + - kind: ServiceAccount + name: {{ include "streaming-operator.serviceAccountName" . }} + namespace: {{ .Release.Namespace }} +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: Role +metadata: + name: {{ include "streaming-operator.fullname" . }} + namespace: {{ include "streaming-operator.workloadNamespace" . }} + labels: + {{- include "streaming-operator.labels" . | nindent 4 }} +rules: + - apiGroups: [""] + resources: [configmaps] + verbs: [get, list, create, patch, delete] + - apiGroups: [apps] + resources: [deployments] + verbs: [get, list, create, patch, delete] +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: RoleBinding +metadata: + name: {{ include "streaming-operator.fullname" . }} + namespace: {{ include "streaming-operator.workloadNamespace" . }} + labels: + {{- include "streaming-operator.labels" . | nindent 4 }} +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: Role + name: {{ include "streaming-operator.fullname" . }} +subjects: + - kind: ServiceAccount + name: {{ include "streaming-operator.serviceAccountName" . }} + namespace: {{ .Release.Namespace }} diff --git a/sentry_streams_k8s/chart/streaming-operator/templates/serviceaccount.yaml b/sentry_streams_k8s/chart/streaming-operator/templates/serviceaccount.yaml new file mode 100644 index 00000000..44cfd13e --- /dev/null +++ b/sentry_streams_k8s/chart/streaming-operator/templates/serviceaccount.yaml @@ -0,0 +1,7 @@ +apiVersion: v1 +kind: ServiceAccount +metadata: + name: {{ include "streaming-operator.serviceAccountName" . }} + namespace: {{ .Release.Namespace }} + labels: + {{- include "streaming-operator.labels" . | nindent 4 }} diff --git a/sentry_streams_k8s/chart/streaming-operator/values.yaml b/sentry_streams_k8s/chart/streaming-operator/values.yaml new file mode 100644 index 00000000..f422782e --- /dev/null +++ b/sentry_streams_k8s/chart/streaming-operator/values.yaml @@ -0,0 +1,32 @@ +# - labels is for things like Sentry COGS tags (service, app_feature, component). +# - serviceAccount.name defaults to the chart's full name when left empty. +# - image.tag defaults to the chart's appVersion when left empty. + +workloadNamespace: "" + +image: + repository: us-central1-docker.pkg.dev/sentryio/streaming-operator/image + tag: "" + pullPolicy: IfNotPresent + +resources: + requests: + cpu: 100m + memory: 256Mi + limits: + memory: 256Mi + +env: [] + +labels: {} + +podLabels: {} +podAnnotations: {} + +serviceAccount: + name: "" + +nodeSelector: {} +tolerations: [] +affinity: {} +priorityClassName: "" diff --git a/sentry_streams_k8s/pyproject.toml b/sentry_streams_k8s/pyproject.toml index 2b3fe3a4..afb25cd6 100644 --- a/sentry_streams_k8s/pyproject.toml +++ b/sentry_streams_k8s/pyproject.toml @@ -50,7 +50,11 @@ where = ["."] include = ["sentry_streams_k8s*"] [tool.setuptools.package-data] -sentry_streams_k8s = ["py.typed", "templates/deployment.yaml", "templates/container.yaml"] +sentry_streams_k8s = [ + "py.typed", + "templates/deployment.yaml", + "templates/container.yaml", +] [tool.pytest.ini_options] minversion = "6.0" diff --git a/sentry_streams_k8s/sentry_streams_k8s/operator/operator.py b/sentry_streams_k8s/sentry_streams_k8s/operator/operator.py index cb4a1554..3f929519 100644 --- a/sentry_streams_k8s/sentry_streams_k8s/operator/operator.py +++ b/sentry_streams_k8s/sentry_streams_k8s/operator/operator.py @@ -1,12 +1,14 @@ from __future__ import annotations import logging +import os from typing import Any import kopf from kubernetes import client, dynamic +from kubernetes.client.exceptions import ApiException -from sentry_streams_k8s.consumer_builder import compute_config_version, make_k8s_name +from sentry_streams_k8s.consumer_builder import compute_config_version from sentry_streams_k8s.operator.streaming_pipeline import ( from_crd_spec, render, @@ -19,37 +21,114 @@ VERSION = "v1alpha1" PLURAL = "streamingpipelines" FIELD_MANAGER = "streaming-operator" +WORKLOAD_NAMESPACE_ENV = "WORKLOAD_NAMESPACE" +OWNER_UID_LABEL = "streams.sentry.io/owner-uid" +OWNER_NAME_ANNOTATION = "streams.sentry.io/owner-name" +OWNER_NAMESPACE_ANNOTATION = "streams.sentry.io/owner-namespace" -def _apply(dyn: dynamic.DynamicClient, manifest: dict[str, Any], namespace: str) -> None: +def _workload_namespace() -> str: + namespace = os.environ.get(WORKLOAD_NAMESPACE_ENV, "").strip() + if not namespace: + raise RuntimeError(f"{WORKLOAD_NAMESPACE_ENV} must be set.") + return namespace + + +def _prepare_manifest( + manifest: dict[str, Any], + *, + workload_namespace: str, + owner_uid: str, + owner_name: str, + owner_namespace: str, +) -> None: + metadata = manifest.setdefault("metadata", {}) + metadata["namespace"] = workload_namespace + metadata["labels"] = { + **metadata.get("labels", {}), + OWNER_UID_LABEL: owner_uid, + } + metadata["annotations"] = { + **metadata.get("annotations", {}), + OWNER_NAME_ANNOTATION: owner_name, + OWNER_NAMESPACE_ANNOTATION: owner_namespace, + } + + +def _apply( + dyn: dynamic.DynamicClient, + manifest: dict[str, Any], + *, + workload_namespace: str, + owner_uid: str, +) -> None: resource = dyn.resources.get(api_version=manifest["apiVersion"], kind=manifest["kind"]) + name = manifest["metadata"]["name"] + try: + existing = resource.get(name=name, namespace=workload_namespace) + except ApiException as e: + if e.status != 404: + raise + else: + labels = existing.metadata.labels or {} + existing_owner_uid = labels.get(OWNER_UID_LABEL) + if existing_owner_uid != owner_uid: + raise kopf.PermanentError( + f"{manifest['kind']} {workload_namespace}/{name} is already present and is not " + "managed by this StreamingPipeline." + ) + dyn.server_side_apply( resource, body=manifest, - namespace=namespace, + namespace=workload_namespace, field_manager=FIELD_MANAGER, force_conflicts=True, ) -def _prune_stale_deployments( +def _prune_stale_resources( *, - namespace: str, + workload_namespace: str, owner_uid: str, - service_name: str, - pipeline_name: str, - desired_names: set[str], + desired_deployments: set[str], + desired_configmaps: set[str], ) -> None: + selector = f"{OWNER_UID_LABEL}={owner_uid}" + apps = client.AppsV1Api() - selector = f"service={make_k8s_name(service_name)},pipeline={make_k8s_name(pipeline_name)}" - existing = apps.list_namespaced_deployment(namespace=namespace, label_selector=selector) - for item in existing.items: - owner_refs = item.metadata.owner_references or [] - if not any(ref.uid == owner_uid for ref in owner_refs): - continue - if item.metadata.name not in desired_names: - logger.info("Pruning stale deployment %s/%s", namespace, item.metadata.name) - apps.delete_namespaced_deployment(name=item.metadata.name, namespace=namespace) + deployments = apps.list_namespaced_deployment( + namespace=workload_namespace, + label_selector=selector, + ) + for deployment in deployments.items: + if deployment.metadata.name not in desired_deployments: + logger.info( + "Pruning stale deployment %s/%s", + workload_namespace, + deployment.metadata.name, + ) + apps.delete_namespaced_deployment( + name=deployment.metadata.name, + namespace=workload_namespace, + ) + + core = client.CoreV1Api() + configmaps = core.list_namespaced_config_map( + namespace=workload_namespace, + label_selector=selector, + ) + for configmap in configmaps.items: + if configmap.metadata.name not in desired_configmaps: + logger.info( + "Pruning stale configmap %s/%s", + workload_namespace, + configmap.metadata.name, + ) + core.delete_namespaced_config_map( + name=configmap.metadata.name, + namespace=workload_namespace, + ) def _condition(type_: str, status: bool, reason: str, message: str = "") -> dict[str, Any]: @@ -73,6 +152,7 @@ def reconcile( **_: Any, ) -> None: assert namespace is not None + workload_namespace = _workload_namespace() consumer = from_crd_spec(dict(spec), name=name) try: @@ -88,15 +168,33 @@ def reconcile( dyn = dynamic.DynamicClient(client.ApiClient()) for manifest in manifests: - kopf.adopt(manifest) - _apply(dyn, manifest, namespace) - - _prune_stale_deployments( - namespace=namespace, + _prepare_manifest( + manifest, + workload_namespace=workload_namespace, + owner_uid=uid, + owner_name=name, + owner_namespace=namespace, + ) + _apply( + dyn, + manifest, + workload_namespace=workload_namespace, + owner_uid=uid, + ) + + _prune_stale_resources( + workload_namespace=workload_namespace, owner_uid=uid, - service_name=consumer["service_name"], - pipeline_name=consumer["pipeline_name"], - desired_names={m["metadata"]["name"] for m in manifests if m["kind"] == "Deployment"}, + desired_deployments={ + manifest["metadata"]["name"] + for manifest in manifests + if manifest["kind"] == "Deployment" + }, + desired_configmaps={ + manifest["metadata"]["name"] + for manifest in manifests + if manifest["kind"] == "ConfigMap" + }, ) replicas = consumer.get("replicas", 1) @@ -107,9 +205,21 @@ def reconcile( ] patch.status["config_version"] = compute_config_version(consumer["pipeline_config"]) patch.status["replicas"] = {"primary": replicas - canary, "canary": canary} + patch.status["workload_namespace"] = workload_namespace + + +@kopf.on.delete(GROUP, VERSION, PLURAL) +def cleanup(uid: str, **_: Any) -> None: + _prune_stale_resources( + workload_namespace=_workload_namespace(), + owner_uid=uid, + desired_deployments=set(), + desired_configmaps=set(), + ) def main() -> None: + _workload_namespace() kopf.run(standalone=True, clusterwide=True) diff --git a/sentry_streams_k8s/tests/test_operator.py b/sentry_streams_k8s/tests/test_operator.py index 91c0d17e..c34859c0 100644 --- a/sentry_streams_k8s/tests/test_operator.py +++ b/sentry_streams_k8s/tests/test_operator.py @@ -1,174 +1,141 @@ from __future__ import annotations from types import SimpleNamespace -from typing import Any -from unittest.mock import MagicMock, call +from unittest.mock import MagicMock, patch import kopf import pytest -from sentry_streams_k8s.consumer_builder import compute_config_version -from sentry_streams_k8s.operator import operator - - -def consumer_spec(**overrides: Any) -> dict[str, Any]: - spec: dict[str, Any] = { - "service_name": "my-service", - "pipeline_name": "my-pipeline", - "pipeline_module": "pipelines/my_pipeline.py", - "image_name": "registry.example.com/image:abc123", - "cpu_per_process": 1000, - "memory_per_process": 512, - "replicas": 4, - "deployment_template": { - "metadata": {"labels": {"component": "my-component"}}, - "spec": { - "selector": {"matchLabels": {"component": "my-component"}}, - "template": { - "metadata": {"labels": {"component": "my-component"}}, - "spec": {"serviceAccountName": "some-service-account"}, - }, - }, - }, - "container_template": {}, - "pipeline_config": { - "env": {}, - "pipeline": { - "segments": [ - { - "steps_config": { - "myinput": { - "starts_segment": True, - "bootstrap_servers": ["127.0.0.1:9092"], - } - } - } - ] - }, +from sentry_streams_k8s.operator.operator import ( + OWNER_NAME_ANNOTATION, + OWNER_NAMESPACE_ANNOTATION, + OWNER_UID_LABEL, + _apply, + _prepare_manifest, + _prune_stale_resources, +) + +WORKLOAD_NAMESPACE = "test-streaming-pipelines" + + +def test_prepare_manifest_routes_workload_and_records_source_cr() -> None: + manifest = { + "apiVersion": "apps/v1", + "kind": "Deployment", + "metadata": { + "name": "pipeline", + "labels": {"service": "test"}, + "annotations": {"existing": "annotation"}, + "namespace": "source", }, } - spec.update(overrides) - return spec - - -def mock_k8s(monkeypatch: pytest.MonkeyPatch) -> tuple[MagicMock, MagicMock, MagicMock]: - api_client = MagicMock(name="api_client") - dynamic_client = MagicMock(name="dynamic_client") - apps_api = MagicMock(name="apps_api") - apps_api.list_namespaced_deployment.return_value = SimpleNamespace(items=[]) - monkeypatch.setattr(operator.client, "ApiClient", MagicMock(return_value=api_client)) - monkeypatch.setattr( - operator.dynamic, "DynamicClient", MagicMock(return_value=dynamic_client) + _prepare_manifest( + manifest, + workload_namespace=WORKLOAD_NAMESPACE, + owner_uid="owner-uid", + owner_name="pipeline-cr", + owner_namespace="default", ) - monkeypatch.setattr(operator.client, "AppsV1Api", MagicMock(return_value=apps_api)) - adopt = MagicMock() - monkeypatch.setattr(operator.kopf, "adopt", adopt) - return dynamic_client, apps_api, adopt + assert manifest["metadata"] == { + "name": "pipeline", + "namespace": WORKLOAD_NAMESPACE, + "labels": { + "service": "test", + OWNER_UID_LABEL: "owner-uid", + }, + "annotations": { + "existing": "annotation", + OWNER_NAME_ANNOTATION: "pipeline-cr", + OWNER_NAMESPACE_ANNOTATION: "default", + }, + } -def test_reconcile_applies_rendered_resources_and_updates_status( - monkeypatch: pytest.MonkeyPatch, -) -> None: - dynamic_client, apps_api, adopt = mock_k8s(monkeypatch) - spec = consumer_spec(with_canary=True) - patch = SimpleNamespace(status={}) - - operator.reconcile( - spec=spec, - name="my-pipeline", - namespace="streams", - uid="pipeline-uid", - patch=patch, - ) - manifests = [args[1]["body"] for args in dynamic_client.server_side_apply.call_args_list] - assert [(manifest["kind"], manifest["metadata"]["name"]) for manifest in manifests] == [ - ("ConfigMap", "my-service-pipeline-my-pipeline"), - ("Deployment", "my-service-pipeline-my-pipeline-0"), - ("Deployment", "my-service-pipeline-my-pipeline-0-canary"), - ] - assert adopt.call_args_list == [call(manifest) for manifest in manifests] - for applied in dynamic_client.server_side_apply.call_args_list: - assert applied.kwargs["namespace"] == "streams" - assert applied.kwargs["field_manager"] == operator.FIELD_MANAGER - assert applied.kwargs["force_conflicts"] is True - - apps_api.list_namespaced_deployment.assert_called_once_with( - namespace="streams", - label_selector="service=my-service,pipeline=my-pipeline", +def test_apply_rejects_resource_owned_by_another_cr() -> None: + resource = MagicMock() + resource.get.return_value = SimpleNamespace( + metadata=SimpleNamespace(labels={OWNER_UID_LABEL: "another-owner"}) ) - assert patch.status == { - "conditions": [ - {"type": "Rendered", "status": "True", "reason": "Rendered", "message": ""}, - {"type": "Applied", "status": "True", "reason": "Applied", "message": ""}, - ], - "config_version": compute_config_version(spec["pipeline_config"]), - "replicas": {"primary": 3, "canary": 1}, + dyn = MagicMock() + dyn.resources.get.return_value = resource + manifest = { + "apiVersion": "v1", + "kind": "ConfigMap", + "metadata": {"name": "pipeline", "namespace": WORKLOAD_NAMESPACE}, } - -def test_reconcile_prunes_only_owned_stale_deployments( - monkeypatch: pytest.MonkeyPatch, -) -> None: - _, apps_api, _ = mock_k8s(monkeypatch) - - def deployment(name: str, owner_uid: str) -> SimpleNamespace: - return SimpleNamespace( - metadata=SimpleNamespace( - name=name, - owner_references=[SimpleNamespace(uid=owner_uid)], - ) + with pytest.raises(kopf.PermanentError, match="not managed by this StreamingPipeline"): + _apply( + dyn, + manifest, + workload_namespace=WORKLOAD_NAMESPACE, + owner_uid="owner-uid", ) - apps_api.list_namespaced_deployment.return_value = SimpleNamespace( - items=[ - deployment("my-service-pipeline-my-pipeline-0", "pipeline-uid"), - deployment("stale-owned", "pipeline-uid"), - deployment("stale-owned-by-another-resource", "other-uid"), - ] + dyn.server_side_apply.assert_not_called() + + +def test_apply_uses_workload_namespace_and_stable_field_manager() -> None: + resource = MagicMock() + resource.get.return_value = SimpleNamespace( + metadata=SimpleNamespace(labels={OWNER_UID_LABEL: "owner-uid"}) ) + dyn = MagicMock() + dyn.resources.get.return_value = resource + manifest = { + "apiVersion": "v1", + "kind": "ConfigMap", + "metadata": {"name": "pipeline", "namespace": WORKLOAD_NAMESPACE}, + } - operator.reconcile( - spec=consumer_spec(), - name="my-pipeline", - namespace="streams", - uid="pipeline-uid", - patch=SimpleNamespace(status={}), + _apply( + dyn, + manifest, + workload_namespace=WORKLOAD_NAMESPACE, + owner_uid="owner-uid", ) - apps_api.delete_namespaced_deployment.assert_called_once_with( - name="stale-owned", - namespace="streams", + resource.get.assert_called_once_with(name="pipeline", namespace=WORKLOAD_NAMESPACE) + dyn.server_side_apply.assert_called_once_with( + resource, + body=manifest, + namespace=WORKLOAD_NAMESPACE, + field_manager="streaming-operator", + force_conflicts=True, ) -def test_reconcile_records_render_failure_without_calling_k8s( - monkeypatch: pytest.MonkeyPatch, +@patch("sentry_streams_k8s.operator.operator.client.CoreV1Api") +@patch("sentry_streams_k8s.operator.operator.client.AppsV1Api") +def test_prune_removes_only_stale_resources( + apps_api: MagicMock, + core_api: MagicMock, ) -> None: - api_client = MagicMock() - monkeypatch.setattr(operator.client, "ApiClient", api_client) - patch = SimpleNamespace(status={}) - - with pytest.raises(kopf.PermanentError, match="failed to render"): - operator.reconcile( - spec={"service_name": "my-service"}, - name="my-pipeline", - namespace="streams", - uid="pipeline-uid", - patch=patch, - ) - - api_client.assert_not_called() - assert patch.status["conditions"] == [ - { - "type": "Rendered", - "status": "False", - "reason": "ValueError", - "message": ( - "StreamingPipeline is missing required field(s): container_template, " - "cpu_per_process, deployment_template, image_name, memory_per_process, " - "pipeline_config, pipeline_module." - ), - } + apps = apps_api.return_value + apps.list_namespaced_deployment.return_value.items = [ + SimpleNamespace(metadata=SimpleNamespace(name="desired-deployment")), + SimpleNamespace(metadata=SimpleNamespace(name="stale-deployment")), + ] + core = core_api.return_value + core.list_namespaced_config_map.return_value.items = [ + SimpleNamespace(metadata=SimpleNamespace(name="desired-configmap")), + SimpleNamespace(metadata=SimpleNamespace(name="stale-configmap")), ] + + _prune_stale_resources( + workload_namespace=WORKLOAD_NAMESPACE, + owner_uid="owner-uid", + desired_deployments={"desired-deployment"}, + desired_configmaps={"desired-configmap"}, + ) + + apps.delete_namespaced_deployment.assert_called_once_with( + name="stale-deployment", + namespace=WORKLOAD_NAMESPACE, + ) + core.delete_namespaced_config_map.assert_called_once_with( + name="stale-configmap", + namespace=WORKLOAD_NAMESPACE, + ) diff --git a/sentry_streams_k8s/tests/test_streaming_pipeline.py b/sentry_streams_k8s/tests/test_streaming_pipeline.py index 36224909..cca4ddbf 100644 --- a/sentry_streams_k8s/tests/test_streaming_pipeline.py +++ b/sentry_streams_k8s/tests/test_streaming_pipeline.py @@ -18,9 +18,9 @@ CRD_PATH = ( pathlib.Path(__file__).resolve().parents[1] - / "sentry_streams_k8s" - / "operator" - / "manifests" + / "chart" + / "streaming-operator" + / "crds" / "crd.yaml" ) @@ -199,10 +199,3 @@ def test_crd_required_fields_match_operator_validation() -> None: version = crd["spec"]["versions"][0] crd_required = version["schema"]["openAPIV3Schema"]["properties"]["spec"]["required"] assert set(crd_required) == set(REQUIRED_FIELDS) - - -def test_library_ships_only_the_crd() -> None: - """Operator deployment manifests (namespace, RBAC, Deployment, ...) are - cluster-specific and belong to the deployment repo, not the library.""" - manifests_dir = CRD_PATH.parent - assert [p.name for p in sorted(manifests_dir.iterdir())] == ["crd.yaml"]