diff --git a/.github/workflows/scripts/before_install.sh b/.github/workflows/scripts/before_install.sh index 91ede8e..c83ed89 100755 --- a/.github/workflows/scripts/before_install.sh +++ b/.github/workflows/scripts/before_install.sh @@ -44,7 +44,7 @@ legacy_component_name: "pulp_workflow" component_name: "workflow" component_version: "${COMPONENT_VERSION}" pulp_env: {} -pulp_settings: {"api_root": "/pulp/"} +pulp_settings: {"WORKFLOW_CALLBACK_FIELDS": ["name", "state", "labels:email"], "api_root": "/pulp/"} pulp_scheme: "https" image: name: "pulp" diff --git a/CHANGES/10.feature b/CHANGES/10.feature new file mode 100644 index 0000000..85386d3 --- /dev/null +++ b/CHANGES/10.feature @@ -0,0 +1,5 @@ +Added user-registered callbacks that fire on `Workflow` lifecycle events. A new `CallbackService` +resource (modeled after pulpcore's `SigningService`) points at an absolute path to an executable; it +is attached to a workflow via a `WorkflowCallback` whose `callback_type` selects the event +(`running`, `completed`, etc). The script runs as a Pulp task with workflow context available as +environment variables. diff --git a/README.md b/README.md index 60bd6e1..e334699 100644 --- a/README.md +++ b/README.md @@ -2,8 +2,12 @@ > **Warning:** This is a community plugin and is not officially supported. Scheduling tasks incorrectly can cause serious issues in your Pulp instance. Always test in a development environment first before applying changes to production. -A Pulp plugin that introduces `Workflow` — a named, ordered pipeline of tasks -dispatched sequentially. +A Pulp plugin that introduces the `Workflow` model. Workflows build on top of +tasks in Pulp allowing users to: +* Schedule tasks to run at any given time +* Run sequences of tasks in a specific order +* Set up callback services to run on workflow lifecycle events (e.g. running, +completed, failed, canceled, finished) A `Workflow` owns one or more `WorkflowTask` rows. Each task records the `task_name`, `task_args`, `task_kwargs`, and any `reserved_resources` to use @@ -14,10 +18,10 @@ workflow, cancel it (if it has not yet started) and create a new one. | Method | URL | Description | |--------|-----|-------------| -| GET | `/pulp/api/v3/workflows/` | List workflows | -| POST | `/pulp/api/v3/workflows/` | Create a workflow (with tasks) | -| GET | `/pulp/api/v3/workflows//` | Retrieve a workflow | -| PATCH | `/pulp/api/v3/workflows//` | Cancel a waiting workflow (body: `{"state": "canceled"}`). Returns 409 if the workflow has already started; only `"canceled"` is accepted as the target state. | +| GET | `/pulp/api/v3/workflow/workflows/` | List workflows | +| POST | `/pulp/api/v3/workflow/workflows/` | Create a workflow (with tasks) | +| GET | `/pulp/api/v3/workflow/workflows//` | Retrieve a workflow | +| PATCH | `/pulp/api/v3/workflow/workflows//` | Cancel a waiting workflow (body: `{"state": "canceled"}`). Returns 409 if the workflow has already started; only `"canceled"` is accepted as the target state. | ## How execution works @@ -110,4 +114,4 @@ the group. Membership means: The group's `all_tasks_dispatched` flag is `False` while the workflow is running and flipped to `True` exactly once the workflow reaches a terminal -state (`completed`, `failed`, or `canceled`). \ No newline at end of file +state (`completed`, `failed`, or `canceled`). diff --git a/docs/demo/README.md b/docs/demo/README.md new file mode 100644 index 0000000..59fbda5 --- /dev/null +++ b/docs/demo/README.md @@ -0,0 +1,132 @@ +# Demo: sync + publish a file repo via a Workflow, with a messaging callback + +Drives `pulp_workflow` end-to-end: register a `CallbackService` (a script that +POSTs to a webhook, e.g. Discord/Slack), create a `file` repo + remote, POST a +`Workflow` that syncs then publishes with a `finished` callback, and watch it +run. + +Assumes a running dev stack (`oci-env compose up`), a configured pulp-cli, and +`NOTIFY_WEBHOOK` available to the Pulp worker. With `oci-env`, add it to +`compose.env` and bounce the stack so the new value is picked up. For testing, +[`httpbin`](https://httpbin.org) returns 200 for any POST: + +```bash +echo 'NOTIFY_WEBHOOK=https://httpbin.org/post' >> ../oci_env/compose.env +oci-env compose up -d # recreate the pulp container so env_file is re-read +``` + +The demo uses fixed resource names, so run `oci-env pdbreset` between runs to +wipe the Pulp DB. + + +## 1. Point at the notify script + +`oci-env` bind-mounts your `pulp-dev/` checkout at `/src` inside the `pulp` +container, so [notify.sh](notify.sh) is visible to the worker at: + +```bash +SCRIPT_PATH=/src/pulp_workflow/docs/demo/notify.sh +oci-env exec ls -la "$SCRIPT_PATH" # sanity check: file is present and +x +``` + +## 2. Register the CallbackService + +```bash +CB_HREF=$(http POST :5001/pulp/api/v3/workflow/callback-services/ \ + name=demo-messaging-notify \ + script=$SCRIPT_PATH \ + | jq -r .pulp_href) +echo "callback_service: $CB_HREF" +``` + +## 3. Create the repository and remote + +```bash +REPO_NAME=demo-file-repo +REMOTE_NAME=demo-file-remote + +pulp file repository create --name "$REPO_NAME" +pulp file remote create \ + --name "$REMOTE_NAME" \ + --url https://fixtures.pulpproject.org/file/PULP_MANIFEST \ + --policy immediate + +REPO_HREF=$(pulp file repository show --name "$REPO_NAME" | jq -r .pulp_href) +REMOTE_HREF=$(pulp file remote show --name "$REMOTE_NAME" | jq -r .pulp_href) + +# pulpcore tasks take pks, not hrefs, so strip the trailing UUID segment. +REPO_PK=$(echo "$REPO_HREF" | awk -F/ '{print $(NF-1)}') +REMOTE_PK=$(echo "$REMOTE_HREF" | awk -F/ '{print $(NF-1)}') +``` + +## 4. Create the workflow + +The `publish` task's `repository_version_pk` is a *dynamic* arg +(`content_type: core.repositoryversion`) — the workflow engine resolves it at +dispatch time from the previous task's `created_resources`. One `finished` +callback fires on any terminal state. + +```bash +WF_HREF=$(http POST :5001/pulp/api/v3/workflow/workflows/ </dev/null diff --git a/pulp_workflow/app/migrations/0004_callbackservice_workflowcallback.py b/pulp_workflow/app/migrations/0004_callbackservice_workflowcallback.py new file mode 100644 index 0000000..c29fa29 --- /dev/null +++ b/pulp_workflow/app/migrations/0004_callbackservice_workflowcallback.py @@ -0,0 +1,51 @@ +# Generated by Django 5.2.13 on 2026-05-08 17:32 + +import django.db.models.deletion +import django_lifecycle.mixins +import pulpcore.app.models.base +import pulpcore.app.util +from django.db import migrations, models + + +class Migration(migrations.Migration): + + dependencies = [ + ('core', '0150_taskschedule_task_kwargs'), + ('workflow', '0003_workflow_task_group'), + ] + + operations = [ + migrations.CreateModel( + name='CallbackService', + fields=[ + ('pulp_id', models.UUIDField(default=pulpcore.app.models.base.pulp_uuid, editable=False, primary_key=True, serialize=False)), + ('pulp_created', models.DateTimeField(auto_now_add=True)), + ('pulp_last_updated', models.DateTimeField(auto_now=True, null=True)), + ('name', models.TextField()), + ('script', models.TextField()), + ('pulp_domain', models.ForeignKey(default=pulpcore.app.util.get_domain_pk, on_delete=django.db.models.deletion.CASCADE, to='core.domain')), + ], + options={ + 'permissions': [('manage_roles_callbackservice', 'Can manage role assignments on callback services')], + 'default_permissions': ('add', 'change', 'delete', 'view'), + 'unique_together': {('pulp_domain', 'name')}, + }, + bases=(django_lifecycle.mixins.LifecycleModelMixin, models.Model), + ), + migrations.CreateModel( + name='WorkflowCallback', + fields=[ + ('pulp_id', models.UUIDField(default=pulpcore.app.models.base.pulp_uuid, editable=False, primary_key=True, serialize=False)), + ('pulp_created', models.DateTimeField(auto_now_add=True)), + ('pulp_last_updated', models.DateTimeField(auto_now=True, null=True)), + ('callback_type', models.TextField(choices=[('running', 'Running'), ('completed', 'Completed'), ('failed', 'Failed'), ('canceled', 'Canceled'), ('finished', 'Finished')])), + ('callback_service', models.ForeignKey(on_delete=django.db.models.deletion.PROTECT, related_name='workflow_callbacks', to='workflow.callbackservice')), + ('dispatched_task', models.ForeignKey(null=True, on_delete=django.db.models.deletion.SET_NULL, related_name='+', to='core.task')), + ('workflow', models.ForeignKey(on_delete=django.db.models.deletion.CASCADE, related_name='callbacks', to='workflow.workflow')), + ], + options={ + 'unique_together': {('workflow', 'callback_service', 'callback_type')}, + }, + bases=(django_lifecycle.mixins.LifecycleModelMixin, models.Model), + ), + ] diff --git a/pulp_workflow/app/models.py b/pulp_workflow/app/models.py index 0924e5f..1c8b2ed 100644 --- a/pulp_workflow/app/models.py +++ b/pulp_workflow/app/models.py @@ -1,12 +1,62 @@ +import asyncio +import json +import os +import re +import subprocess +from gettext import gettext as _ + +from django.conf import settings from django.contrib.postgres.fields import ArrayField, HStoreField +from django.core.exceptions import ImproperlyConfigured, ValidationError from django.db import models from django.utils import timezone +from django_guid import get_guid from pulpcore.plugin.constants import TASK_CHOICES, TASK_STATES from pulpcore.plugin.models import BaseModel, EncryptedJSONField from pulpcore.plugin.util import get_domain_pk +# --------------------------------------------------------------------------- +# Callback type constants. These are the events on which a CallbackService can be triggered. +# The first set mirror Workflow lifecycle states; ``FINISHED`` is a synthetic type that fires on +# any terminal state (completed, failed, canceled). +# --------------------------------------------------------------------------- +class CALLBACK_TYPES: # noqa: N801 - mirror pulpcore.constants style (TASK_STATES, ...) + RUNNING = TASK_STATES.RUNNING + COMPLETED = TASK_STATES.COMPLETED + FAILED = TASK_STATES.FAILED + CANCELED = TASK_STATES.CANCELED + # Wildcard: fires on any terminal state. + FINISHED = "finished" + + +CALLBACK_TYPE_CHOICES = ( + (CALLBACK_TYPES.RUNNING, "Running"), + (CALLBACK_TYPES.COMPLETED, "Completed"), + (CALLBACK_TYPES.FAILED, "Failed"), + (CALLBACK_TYPES.CANCELED, "Canceled"), + (CALLBACK_TYPES.FINISHED, "Finished"), +) + +# Map a workflow state transition to the set of callback types that should fire. The key is the +# workflow's new state; the value is the tuple of CallbackService callback_type values to dispatch. +TRANSITION_CALLBACK_TYPES = { + TASK_STATES.RUNNING: (CALLBACK_TYPES.RUNNING,), + TASK_STATES.COMPLETED: (CALLBACK_TYPES.COMPLETED, CALLBACK_TYPES.FINISHED), + TASK_STATES.FAILED: (CALLBACK_TYPES.FAILED, CALLBACK_TYPES.FINISHED), + TASK_STATES.CANCELED: (CALLBACK_TYPES.CANCELED, CALLBACK_TYPES.FINISHED), +} + +# Env var keys must be POSIX-portable: [A-Z_][A-Z0-9_]*. Sanitize label keys to fit that shape +# before exposing them as PULP_WORKFLOW_LABEL_. +_ENV_KEY_RE = re.compile(r"[^A-Z0-9_]") + +# Valid scalar entries for ``WORKFLOW_CALLBACK_FIELDS``. ``labels:`` is also accepted to +# expose a single label without leaking the rest. +ALLOWED_CALLBACK_FIELDS = frozenset({"pk", "name", "state", "labels"}) + + class Workflow(BaseModel): """ A named, ordered pipeline of tasks executed sequentially. @@ -173,3 +223,190 @@ class Meta(_WorkflowTaskArgBase.Meta): name="workflowtaskkwarg_value_ctype_exclusive", ), ] + + +class CallbackService(BaseModel): + """ + A user-registered subprocess invoked when a Workflow reaches a lifecycle event. + + Modeled after pulpcore's ``SigningService``: the ``script`` field is an absolute path to an + executable on the Pulp worker host. At registration time the script is validated for existence + and the executable bit; when invoked it is run as a subprocess with workflow context exposed + via environment variables (``PULP_WORKFLOW_NAME``, ``PULP_WORKFLOW_STATE``, etc.). Which + workflow fields are exposed is controlled by the ``WORKFLOW_CALLBACK_FIELDS`` setting. + + Unlike ``SigningService`` (which is admin-installed out of band), ``CallbackService`` is fully + API-managed and RBAC-scoped. + """ + + name = models.TextField() + script = models.TextField() + + pulp_domain = models.ForeignKey("core.Domain", default=get_domain_pk, on_delete=models.CASCADE) + + def __str__(self): + return f"CallbackService: {self.name}" + + def _env(self, workflow, env_vars=None): + """Build the env dict passed to the script. Honors ``WORKFLOW_CALLBACK_FIELDS``.""" + guid = get_guid() + env = {"CORRELATION_ID": guid if guid else ""} + + fields = getattr(settings, "WORKFLOW_CALLBACK_FIELDS", ["name", "state"]) or [] + scalar_fields = set() + label_keys = set() + expose_all_labels = False + unknown = [] + for entry in fields: + if entry.startswith("labels:"): + key = entry[len("labels:") :] + if key: + label_keys.add(key) + else: + unknown.append(entry) + elif entry == "labels": + expose_all_labels = True + elif entry in ALLOWED_CALLBACK_FIELDS: + scalar_fields.add(entry) + else: + unknown.append(entry) + if unknown: + raise ImproperlyConfigured( + _("WORKFLOW_CALLBACK_FIELDS contains unknown entries: {unknown!r}.").format( + unknown=sorted(unknown) + ) + ) + + if "pk" in scalar_fields: + env["PULP_WORKFLOW_PK"] = str(workflow.pk) + if "name" in scalar_fields: + env["PULP_WORKFLOW_NAME"] = workflow.name + if "state" in scalar_fields: + env["PULP_WORKFLOW_STATE"] = workflow.state + if expose_all_labels or label_keys: + labels = workflow.pulp_labels or {} + if expose_all_labels: + env["PULP_WORKFLOW_LABELS"] = json.dumps(labels, sort_keys=True) + items = labels.items() + else: + items = ((key, labels.get(key)) for key in label_keys) + for key, value in items: + safe_key = _ENV_KEY_RE.sub("_", key.upper()) + env[f"PULP_WORKFLOW_LABEL_{safe_key}"] = "" if value is None else str(value) + + if env_vars: + env.update(env_vars) + # Inherit the worker's environment so PATH, etc. work in the script. + return {**os.environ, **env} + + def run(self, workflow, env_vars=None): + """Run the script synchronously with workflow context. + + Returns a dict with ``returncode``, ``stdout`` and ``stderr``. Raises ``RuntimeError`` on a + non-zero exit so the surrounding task records the failure. + """ + completed = subprocess.run( + [self.script], + env=self._env(workflow, env_vars=env_vars), + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + ) + result = { + "returncode": completed.returncode, + "stdout": completed.stdout.decode("utf-8", errors="replace"), + "stderr": completed.stderr.decode("utf-8", errors="replace"), + } + if completed.returncode != 0: + raise RuntimeError( + _("CallbackService {name!r} exited with {rc}: {err}").format( + name=self.name, rc=completed.returncode, err=result["stderr"] + ) + ) + return result + + async def arun(self, workflow, env_vars=None): + """Async equivalent of :meth:`run`.""" + process = await asyncio.create_subprocess_exec( + self.script, + env=self._env(workflow, env_vars=env_vars), + stdout=asyncio.subprocess.PIPE, + stderr=asyncio.subprocess.PIPE, + ) + stdout, stderr = await process.communicate() + result = { + "returncode": process.returncode, + "stdout": stdout.decode("utf-8", errors="replace"), + "stderr": stderr.decode("utf-8", errors="replace"), + } + if process.returncode != 0: + raise RuntimeError( + _("CallbackService {name!r} exited with {rc}: {err}").format( + name=self.name, rc=process.returncode, err=result["stderr"] + ) + ) + return result + + def validate(self): + """ + Validate that the script is an absolute path to an executable file. + + Raises ``django.core.exceptions.ValidationError`` if the script cannot be invoked. Called + from :meth:`save` so misconfigured services are rejected before they're persisted, and + from the serializer's ``validate_script`` so API clients get a 400 rather than a 500. + """ + if not self.script: + raise ValidationError(_("`script` is required.")) + if not os.path.isabs(self.script): + raise ValidationError( + _("`script` must be an absolute path, got {p!r}.").format(p=self.script) + ) + if not os.path.isfile(self.script): + raise ValidationError( + _("`script` does not exist or is not a file: {p!r}.").format(p=self.script) + ) + if not os.access(self.script, os.X_OK): + raise ValidationError(_("`script` is not executable: {p!r}.").format(p=self.script)) + + def save(self, *args, **kwargs): + self.validate() + super().save(*args, **kwargs) + + class Meta: + default_permissions = ("add", "change", "delete", "view") + permissions = [ + ("manage_roles_callbackservice", "Can manage role assignments on callback services"), + ] + unique_together = ("pulp_domain", "name") + + +class WorkflowCallback(BaseModel): + """ + Attaches a ``CallbackService`` to a ``Workflow`` for a specific lifecycle event. + + A workflow may have any number of callbacks, but each + ``(workflow, callback_service, callback_type)`` triple is unique. When the workflow reaches the + event named by ``callback_type``, a Pulp task is dispatched that runs the service; the + dispatched task is recorded on ``dispatched_task`` so callers can inspect its result via the + API. + """ + + workflow = models.ForeignKey(Workflow, related_name="callbacks", on_delete=models.CASCADE) + callback_service = models.ForeignKey( + CallbackService, related_name="workflow_callbacks", on_delete=models.PROTECT + ) + callback_type = models.TextField(choices=CALLBACK_TYPE_CHOICES) + dispatched_task = models.ForeignKey( + "core.Task", + null=True, + related_name="+", + on_delete=models.SET_NULL, + ) + + def __str__(self): + return ( + f"WorkflowCallback: {self.workflow.name} -> " + f"{self.callback_service.name} on {self.callback_type}" + ) + + class Meta: + unique_together = ("workflow", "callback_service", "callback_type") diff --git a/pulp_workflow/app/serializers.py b/pulp_workflow/app/serializers.py index 76a2957..96a7e1c 100644 --- a/pulp_workflow/app/serializers.py +++ b/pulp_workflow/app/serializers.py @@ -1,6 +1,7 @@ from gettext import gettext as _ from django.contrib.contenttypes.models import ContentType +from django.core.exceptions import ValidationError as DjangoValidationError from django.db import transaction from rest_framework import serializers from rest_framework.validators import UniqueValidator @@ -8,6 +9,7 @@ from pulpcore.plugin.constants import TASK_CHOICES, TASK_STATES from pulpcore.plugin.models import TaskGroup, TaskSchedule from pulpcore.plugin.serializers import ( + DomainUniqueValidator, IdentityField, ModelSerializer, RelatedField, @@ -15,7 +17,10 @@ ) from pulp_workflow.app.models import ( + CALLBACK_TYPE_CHOICES, + CallbackService, Workflow, + WorkflowCallback, WorkflowTask, WorkflowTaskArg, WorkflowTaskKwarg, @@ -163,10 +168,52 @@ def validate_task_kwargs(self, value): return value +class CallbackServiceRelatedField(RelatedField): + """A hyperlinked relation to a ``CallbackService`` by its detail URL or PRN.""" + + view_name = "workflow-callback-services-detail" + + # ``queryset`` is set in ``__init__`` rather than as a class attribute so importing this module + # does not require Django app loading to be far enough along for ``CallbackService.objects`` to + # resolve. + def __init__(self, **kwargs): + kwargs.setdefault("queryset", CallbackService.objects.all()) + super().__init__(**kwargs) + + +class WorkflowCallbackSerializer(serializers.ModelSerializer): + """A ``WorkflowCallback`` nested under a ``Workflow``. + + On create, ``callback_service`` and ``callback_type`` are required; ``dispatched_task`` is + read-only and is populated by ``execute_workflow`` when the corresponding lifecycle event + fires. + """ + + callback_service = CallbackServiceRelatedField( + help_text=_("Href of the CallbackService to invoke."), + ) + callback_type = serializers.ChoiceField( + choices=CALLBACK_TYPE_CHOICES, + help_text=_( + "The workflow lifecycle event that triggers this callback. The 'finished' " + "type fires on any terminal state (completed, failed, canceled)." + ), + ) + dispatched_task = RelatedField( + view_name="tasks-detail", + read_only=True, + help_text=_("Href of the most recently dispatched callback task, if any."), + ) + + class Meta: + model = WorkflowCallback + fields = ("callback_service", "callback_type", "dispatched_task") + + class WorkflowSerializer(ModelSerializer): """Serializer for Workflow with nested tasks.""" - pulp_href = IdentityField(view_name="workflows-detail") + pulp_href = IdentityField(view_name="workflow-workflows-detail") name = serializers.CharField( help_text=_("The name of the workflow."), allow_blank=False, @@ -227,6 +274,11 @@ class WorkflowSerializer(ModelSerializer): allow_empty=False, help_text=_("The ordered tasks that make up this workflow."), ) + callbacks = WorkflowCallbackSerializer( + many=True, + required=False, + help_text=_("User-registered callbacks that fire on this workflow's lifecycle events."), + ) class Meta: model = Workflow @@ -241,8 +293,22 @@ class Meta: "current_task", "task_group", "tasks", + "callbacks", ) + def validate_callbacks(self, value): + seen = set() + for row in value: + key = (row["callback_service"].pk, row["callback_type"]) + if key in seen: + raise serializers.ValidationError( + _("Duplicate (callback_service, callback_type): ({s!r}, {t!r}).").format( + s=row["callback_service"].name, t=row["callback_type"] + ) + ) + seen.add(key) + return value + def validate_tasks(self, value): # Dynamic args reference the previous task's created resources, so the first task # cannot use them. @@ -257,6 +323,7 @@ def validate_tasks(self, value): @transaction.atomic def create(self, validated_data): tasks_data = validated_data.pop("tasks") + callbacks_data = validated_data.pop("callbacks", []) workflow = Workflow.objects.create(**validated_data) workflow.task_group = TaskGroup.objects.create( description=f"Workflow: {workflow.name}", @@ -274,6 +341,9 @@ def create(self, validated_data): WorkflowTaskKwarg.objects.bulk_create( WorkflowTaskKwarg(workflow_task=wf_task, **row) for row in task_kwargs ) + WorkflowCallback.objects.bulk_create( + WorkflowCallback(workflow=workflow, **row) for row in callbacks_data + ) # Schedule a one-shot dispatch of execute_workflow at start_time. # dispatch_interval=None makes pulpcore's scheduler fire it once and stop. @@ -298,3 +368,36 @@ class WorkflowCancelSerializer(serializers.Serializer): class Meta: fields = ("state",) + + +class CallbackServiceSerializer(ModelSerializer): + """Serializer for ``CallbackService``.""" + + pulp_href = IdentityField(view_name="workflow-callback-services-detail") + name = serializers.CharField( + help_text=_("A name for this callback service. Unique within a domain."), + validators=[DomainUniqueValidator(queryset=CallbackService.objects.all())], + ) + script = serializers.CharField( + help_text=_( + "An absolute path on the Pulp worker host to an executable script that is " + "invoked when an attached workflow reaches the registered callback type. " + "Workflow context is exposed via PULP_WORKFLOW_* environment variables; the " + "exact subset is controlled by the ``WORKFLOW_CALLBACK_FIELDS`` server setting " + "(defaults to PULP_WORKFLOW_NAME and PULP_WORKFLOW_STATE). PULP_WORKFLOW_PK, " + "PULP_WORKFLOW_LABELS, and PULP_WORKFLOW_LABEL_ may also be exposed." + ), + ) + + def validate_script(self, value): + # Run the same checks that ``CallbackService.validate`` enforces on save, but raise a DRF + # ``ValidationError`` so the API returns 400 rather than 500 for misconfigured input. + try: + CallbackService(script=value).validate() + except DjangoValidationError as exc: + raise serializers.ValidationError(exc.messages) + return value + + class Meta: + model = CallbackService + fields = ModelSerializer.Meta.fields + ("name", "script") diff --git a/pulp_workflow/app/settings.py b/pulp_workflow/app/settings.py index 46b71a0..3ec0af7 100644 --- a/pulp_workflow/app/settings.py +++ b/pulp_workflow/app/settings.py @@ -5,3 +5,9 @@ # str/int). Pinning back to OAS 3.0.1 mirrors what pulp_rpm does and keeps the # generated bindings working. SPECTACULAR_SETTINGS__OAS_VERSION = "3.0.1" + +# Workflow fields exposed to ``CallbackService`` scripts via ``PULP_WORKFLOW_*`` env vars. +# Allowed: "pk", "name", "state", "labels" (all labels), or "labels:" (one label). +# ``CORRELATION_ID`` is always exposed. +# Only expose more fields if you're certain it won't present a security risk. +WORKFLOW_CALLBACK_FIELDS = ["name", "state"] diff --git a/pulp_workflow/app/tasks.py b/pulp_workflow/app/tasks.py index 3fcb671..cffc086 100644 --- a/pulp_workflow/app/tasks.py +++ b/pulp_workflow/app/tasks.py @@ -7,7 +7,12 @@ from pulpcore.plugin.constants import TASK_STATES from pulpcore.plugin.tasking import dispatch -from pulp_workflow.app.models import Workflow, WorkflowTask +from pulp_workflow.app.models import ( + TRANSITION_CALLBACK_TYPES, + Workflow, + WorkflowCallback, + WorkflowTask, +) _log = logging.getLogger(__name__) @@ -17,6 +22,54 @@ def _workflow_resource(workflow_pk): return f"pulp_workflow:workflow:{workflow_pk}" +def dispatch_workflow_callbacks(workflow, new_state): + """Dispatch a ``run_callback`` task for each callback matching ``new_state``. + + Called whenever a Workflow transitions to a new state. The mapping from state to callback types + lives in ``TRANSITION_CALLBACK_TYPES``; states not in that map (e.g. ``waiting``) do not fire + callbacks. + + Best-effort: a failure to dispatch any one callback is logged but does not bubble up so that + one bad callback cannot break the workflow's own state transition. + """ + types = TRANSITION_CALLBACK_TYPES.get(new_state) + if not types: + return + callbacks = workflow.callbacks.filter(callback_type__in=types).select_related( + "callback_service" + ) + for wfcb in callbacks: + try: + child = dispatch( + run_callback, + kwargs={"workflow_callback_pk": str(wfcb.pk)}, + task_group=workflow.task_group, + ) + except Exception: + _log.exception( + "Failed to dispatch callback %s for workflow %s", + wfcb.callback_service.name, + workflow.name, + ) + continue + wfcb.dispatched_task = child + wfcb.save(update_fields=["dispatched_task", "pulp_last_updated"]) + + +def run_callback(workflow_callback_pk): + """Pulp task that invokes a ``CallbackService`` for one ``WorkflowCallback``. + + The callback service runs as a subprocess on the worker host with the workflow's name, state, + pk, and labels exposed via ``PULP_WORKFLOW_*`` environment variables. Non-zero exit raises + ``RuntimeError``, which marks this task as failed (visible via the WorkflowCallback's + ``dispatched_task`` href on the Workflow detail endpoint). + """ + wfcb = WorkflowCallback.objects.select_related("workflow", "callback_service").get( + pk=workflow_callback_pk + ) + return wfcb.callback_service.run(wfcb.workflow) + + def execute_workflow(workflow_pk, next_index=0): """ Run one step of a Workflow, then re-dispatch ourselves for the next step. @@ -50,6 +103,7 @@ def execute_workflow(workflow_pk, next_index=0): workflow.started_at = timezone.now() workflow.save(update_fields=["state", "started_at", "pulp_last_updated"]) _log.info("Workflow %s started.", workflow.name) + dispatch_workflow_callbacks(workflow, TASK_STATES.RUNNING) prev_task = None else: # Continuation: inspect the previous step's child task. @@ -90,6 +144,7 @@ def execute_workflow(workflow_pk, next_index=0): workflow.save(update_fields=["state", "finished_at", "current_task", "pulp_last_updated"]) _mark_task_group_dispatched(workflow) _log.info("Workflow %s completed.", workflow.name) + dispatch_workflow_callbacks(workflow, TASK_STATES.COMPLETED) return workflow.current_task = wf_task @@ -149,6 +204,7 @@ def _fail_workflow(workflow, wf_task, exc=None, description=None, child_error=No _log.info( "Workflow %s failed at step %d (%s).", workflow.name, wf_task.index, wf_task.task_name ) + dispatch_workflow_callbacks(workflow, TASK_STATES.FAILED) def _mark_task_group_dispatched(workflow): diff --git a/pulp_workflow/app/viewsets.py b/pulp_workflow/app/viewsets.py index 653ca40..f8395fc 100644 --- a/pulp_workflow/app/viewsets.py +++ b/pulp_workflow/app/viewsets.py @@ -7,6 +7,7 @@ from pulpcore.plugin.constants import TASK_STATES from pulpcore.plugin.models import TaskSchedule from pulpcore.plugin.viewsets import ( + DATETIME_FILTER_OPTIONS, NAME_FILTER_OPTIONS, BaseFilterSet, LabelFilter, @@ -15,11 +16,32 @@ RolesMixin, ) -from pulp_workflow.app.models import Workflow -from pulp_workflow.app.serializers import WorkflowCancelSerializer, WorkflowSerializer +from pulp_workflow.app.models import CallbackService, Workflow +from pulp_workflow.app.serializers import ( + CallbackServiceSerializer, + WorkflowCancelSerializer, + WorkflowSerializer, +) +from pulp_workflow.app.tasks import dispatch_workflow_callbacks + + +class WorkflowPluginViewSetMixin: + """Mixin that mounts every ``pulp_workflow`` endpoint under ``/workflow/``. + + Pulpcore does not automatically scope plain (non-Master/Detail) plugin viewsets under their + plugin name, so ``endpoint_pieces`` is overridden here to prepend ``"workflow"``. This keeps + ``pulp_workflow``'s endpoints (``/pulp/api/v3/workflow/workflows/``, + ``/pulp/api/v3/workflow/callback-services/``) grouped under a stable plugin prefix. -# DATETIME_FILTER_OPTIONS is not re-exported through pulpcore.plugin, so define locally. -DATETIME_FILTER_OPTIONS = ["exact", "lt", "lte", "gt", "gte", "range", "isnull"] + Implemented as a plain mixin rather than a ``NamedModelViewSet`` subclass so that pulpcore's + ``import_viewsets`` does not try to introspect a ``queryset`` on it during app startup. + Must be listed *before* ``NamedModelViewSet`` in a viewset's MRO so ``super()`` resolves to + the right ``endpoint_pieces`` implementation. + """ + + @classmethod + def endpoint_pieces(cls): + return ["workflow", *super().endpoint_pieces()] class WorkflowFilter(BaseFilterSet): @@ -44,6 +66,7 @@ class Meta: class WorkflowViewSet( + WorkflowPluginViewSetMixin, NamedModelViewSet, mixins.CreateModelMixin, mixins.RetrieveModelMixin, @@ -61,9 +84,15 @@ class WorkflowViewSet( queryset = ( Workflow.objects.all() .select_related("task_group", "current_task__dispatched_task") - .prefetch_related("tasks") + .prefetch_related( + "tasks", + "callbacks", + "callbacks__callback_service", + "callbacks__dispatched_task", + ) ) endpoint_name = "workflows" + pulp_tag_name = "Workflows" serializer_class = WorkflowSerializer filterset_class = WorkflowFilter ordering = "-pulp_created" @@ -123,6 +152,7 @@ def partial_update(self, request, pk=None, partial=True): serializer.is_valid(raise_exception=True) workflow = self.get_object() + fired_callbacks = False with transaction.atomic(): workflow = Workflow.objects.select_for_update().get(pk=workflow.pk) if workflow.state == TASK_STATES.WAITING: @@ -135,9 +165,101 @@ def partial_update(self, request, pk=None, partial=True): workflow.task_group.save( update_fields=["all_tasks_dispatched", "pulp_last_updated"] ) + fired_callbacks = True http_status = None else: http_status = status.HTTP_409_CONFLICT + if fired_callbacks: + # Fire CANCELED + FINISHED callbacks. Done outside the select_for_update block so the + # dispatch can read the workflow's persisted state. + dispatch_workflow_callbacks(workflow, TASK_STATES.CANCELED) + out = WorkflowSerializer(workflow, context={"request": request}) return Response(out.data, status=http_status) + + +class CallbackServiceFilter(BaseFilterSet): + """Filter for CallbackServices.""" + + class Meta: + model = CallbackService + fields = { + "name": NAME_FILTER_OPTIONS, + "pulp_created": DATETIME_FILTER_OPTIONS, + } + + +class CallbackServiceViewSet( + WorkflowPluginViewSetMixin, + NamedModelViewSet, + mixins.CreateModelMixin, + mixins.RetrieveModelMixin, + mixins.ListModelMixin, + mixins.UpdateModelMixin, + mixins.DestroyModelMixin, + RolesMixin, +): + """A ViewSet for managing CallbackServices. + + A CallbackService points at an absolute path to an executable on the Pulp worker host. It is + invoked when an attached Workflow reaches a registered lifecycle event. Because callbacks run + arbitrary host scripts they are treated as a privileged resource and require + ``callbackservice_admin``-level permissions to manage. + """ + + queryset = CallbackService.objects.all() + endpoint_name = "callback-services" + pulp_tag_name = "Callback Services" + serializer_class = CallbackServiceSerializer + filterset_class = CallbackServiceFilter + ordering = "-pulp_created" + queryset_filtering_required_permission = "workflow.view_callbackservice" + + DEFAULT_ACCESS_POLICY = { + "statements": [ + { + "action": ["list", "retrieve", "my_permissions"], + "principal": "authenticated", + "effect": "allow", + "condition": "has_model_or_domain_or_obj_perms:workflow.view_callbackservice", + }, + { + "action": ["create"], + "principal": "authenticated", + "effect": "allow", + "condition": "has_model_or_domain_perms:workflow.add_callbackservice", + }, + { + "action": ["update", "partial_update"], + "principal": "authenticated", + "effect": "allow", + "condition": "has_model_or_domain_or_obj_perms:workflow.change_callbackservice", + }, + { + "action": ["destroy"], + "principal": "authenticated", + "effect": "allow", + "condition": "has_model_or_domain_or_obj_perms:workflow.delete_callbackservice", + }, + { + "action": ["list_roles", "add_role", "remove_role"], + "principal": "authenticated", + "effect": "allow", + "condition": ( + "has_model_or_domain_or_obj_perms:workflow.manage_roles_callbackservice" + ), + }, + ], + "queryset_scoping": {"function": "scope_queryset"}, + } + LOCKED_ROLES = { + "workflow.callbackservice_admin": [ + "workflow.view_callbackservice", + "workflow.add_callbackservice", + "workflow.change_callbackservice", + "workflow.delete_callbackservice", + "workflow.manage_roles_callbackservice", + ], + "workflow.callbackservice_viewer": ["workflow.view_callbackservice"], + } diff --git a/pulp_workflow/pytest_plugin.py b/pulp_workflow/pytest_plugin.py index fb647a1..97a111a 100644 --- a/pulp_workflow/pytest_plugin.py +++ b/pulp_workflow/pytest_plugin.py @@ -1,10 +1,38 @@ import uuid from contextlib import suppress +from time import sleep import pytest from pulpcore.tests.functional.utils import BindingsNamespace +# Mirrors the constants pulpcore's monitor_task fixture uses. +WORKFLOW_TIMEOUT = 30 * 60 +WORKFLOW_SLEEP_TIME = 2.0 +WORKFLOW_FINAL_STATES = {"completed", "failed", "canceled", "skipped"} + + +class WorkflowError(Exception): + """Raised when a Workflow reaches a non-completed final state.""" + + def __init__(self, workflow): + self.workflow = workflow + super().__init__( + f"Workflow {workflow.pulp_href} ended in state " + f"{workflow.state!r}: error={workflow.error!r}" + ) + + +class WorkflowTimeoutError(Exception): + """Raised when a Workflow does not reach a final state in the timeout.""" + + def __init__(self, workflow): + self.workflow = workflow + super().__init__( + f"Workflow {workflow.pulp_href} did not reach a final state in time " + f"(state={workflow.state!r})" + ) + @pytest.fixture(scope="session") def workflow_bindings(_api_client_set, bindings_cfg): @@ -60,3 +88,55 @@ def _create_workflow(**kwargs): for href in reversed(created): with suppress(Exception): workflow_bindings.WorkflowsApi.partial_update(href, {"state": "canceled"}) + + +@pytest.fixture +def callback_service_factory(workflow_bindings): + """A factory to generate a CallbackService with auto-cleanup at teardown. + + By default the script is ``/bin/echo`` so the callback always succeeds and is safe to invoke + many times. Override with a ``script`` kwarg to point at a different absolute path. + """ + + created = [] + + def _create_callback_service(**kwargs): + kwargs.setdefault("name", str(uuid.uuid4())) + kwargs.setdefault("script", "/bin/echo") + cs = workflow_bindings.CallbackServicesApi.create(kwargs) + created.append(cs.pulp_href) + return cs + + yield _create_callback_service + + for href in reversed(created): + with suppress(Exception): + workflow_bindings.CallbackServicesApi.delete(href) + + +@pytest.fixture(scope="session") +def monitor_workflow(workflow_bindings): + """Wait for a Workflow to reach a final state. + + Mirrors pulpcore's ``monitor_task`` fixture: returns the Workflow in ``completed`` state, + raises ``WorkflowTimeoutError`` if the timeout in seconds (defaulting to 30*60) is exceeded, + or raises ``WorkflowError`` if it reached any other final state. + """ + + def _monitor_workflow(workflow_href, timeout=WORKFLOW_TIMEOUT): + # Always make at least one read attempt, even if the timeout is shorter than the sleep + # interval, so ``workflow`` is bound before the ``else`` branch can reference it. + attempts = max(1, int(timeout / WORKFLOW_SLEEP_TIME)) + for _ in range(attempts): + workflow = workflow_bindings.WorkflowsApi.read(workflow_href) + if workflow.state in WORKFLOW_FINAL_STATES: + break + sleep(WORKFLOW_SLEEP_TIME) + else: + raise WorkflowTimeoutError(workflow) + + if workflow.state != "completed": + raise WorkflowError(workflow) + return workflow + + return _monitor_workflow diff --git a/pulp_workflow/tests/functional/api/test_crud_callback_services.py b/pulp_workflow/tests/functional/api/test_crud_callback_services.py new file mode 100644 index 0000000..f7d66de --- /dev/null +++ b/pulp_workflow/tests/functional/api/test_crud_callback_services.py @@ -0,0 +1,72 @@ +"""CRUD tests for the CallbackService endpoint.""" + +import uuid + +import pytest + + +@pytest.mark.parallel +def test_create_callback_service(workflow_bindings, callback_service_factory): + """A CallbackService can be created with a name and absolute script path.""" + name = str(uuid.uuid4()) + cs = callback_service_factory(name=name, script="/bin/echo") + assert cs.pulp_href is not None + assert cs.name == name + assert cs.script == "/bin/echo" + + +@pytest.mark.parallel +def test_create_callback_service_rejects_relative_script(workflow_bindings): + """The script must be an absolute path.""" + with pytest.raises(workflow_bindings.ApiException) as exc: + workflow_bindings.CallbackServicesApi.create({"name": str(uuid.uuid4()), "script": "echo"}) + assert exc.value.status == 400 + + +@pytest.mark.parallel +def test_create_callback_service_rejects_missing_script(workflow_bindings): + """The script must point at an existing executable.""" + with pytest.raises(workflow_bindings.ApiException) as exc: + workflow_bindings.CallbackServicesApi.create( + {"name": str(uuid.uuid4()), "script": "/nonexistent/path/to/script"} + ) + assert exc.value.status == 400 + + +@pytest.mark.parallel +def test_create_duplicate_callback_service_name_fails(workflow_bindings, callback_service_factory): + """Names are unique.""" + name = str(uuid.uuid4()) + callback_service_factory(name=name) + with pytest.raises(workflow_bindings.ApiException) as exc: + workflow_bindings.CallbackServicesApi.create({"name": name, "script": "/bin/echo"}) + assert exc.value.status == 400 + + +@pytest.mark.parallel +def test_read_callback_service(workflow_bindings, callback_service_factory): + """A created CallbackService can be retrieved by href.""" + cs = callback_service_factory() + fetched = workflow_bindings.CallbackServicesApi.read(cs.pulp_href) + assert fetched.pulp_href == cs.pulp_href + assert fetched.name == cs.name + + +@pytest.mark.parallel +def test_list_callback_services(workflow_bindings, callback_service_factory): + """Listing and filtering CallbackServices.""" + name = str(uuid.uuid4()) + callback_service_factory(name=name) + results = workflow_bindings.CallbackServicesApi.list(name=name) + assert results.count == 1 + assert results.results[0].name == name + + +@pytest.mark.parallel +def test_delete_callback_service(workflow_bindings, callback_service_factory): + """An unattached CallbackService can be deleted.""" + cs = callback_service_factory() + workflow_bindings.CallbackServicesApi.delete(cs.pulp_href) + with pytest.raises(workflow_bindings.ApiException) as exc: + workflow_bindings.CallbackServicesApi.read(cs.pulp_href) + assert exc.value.status == 404 diff --git a/pulp_workflow/tests/functional/api/test_execute_workflow.py b/pulp_workflow/tests/functional/api/test_execute_workflow.py index 194afb5..bed671f 100644 --- a/pulp_workflow/tests/functional/api/test_execute_workflow.py +++ b/pulp_workflow/tests/functional/api/test_execute_workflow.py @@ -7,27 +7,9 @@ the unique ``RepositoryVersion`` created by task 0. """ -import time import uuid -WORKFLOW_TIMEOUT_SECONDS = 300 -POLL_INTERVAL_SECONDS = 2.0 - - -def _pk_from_href(href): - return href.rstrip("/").split("/")[-1] - - -def _wait_for_workflow(api, workflow_href, timeout=WORKFLOW_TIMEOUT_SECONDS): - """Poll a Workflow until it reaches a final state.""" - final_states = {"completed", "failed", "canceled", "skipped"} - deadline = time.monotonic() + timeout - while time.monotonic() < deadline: - workflow = api.read(workflow_href) - if workflow.state in final_states: - return workflow - time.sleep(POLL_INTERVAL_SECONDS) - raise AssertionError(f"Workflow {workflow_href} did not finish within {timeout}s") +from pulpcore.plugin.util import extract_pk def test_execute_workflow_add_content_and_publish( @@ -37,6 +19,7 @@ def test_execute_workflow_add_content_and_publish( file_repo, file_content_unit_with_name_factory, workflow_factory, + monitor_workflow, ): """A Workflow that adds content then publishes the new version end-to-end.""" repo = file_repo @@ -49,11 +32,11 @@ def test_execute_workflow_add_content_and_publish( "task_kwargs": [ { "kwarg_key": "repository_pk", - "value": _pk_from_href(repo.pulp_href), + "value": extract_pk(repo.pulp_href), }, { "kwarg_key": "add_content_units", - "value": [_pk_from_href(content_a.pulp_href)], + "value": [extract_pk(content_a.pulp_href)], }, {"kwarg_key": "remove_content_units", "value": []}, ], @@ -75,12 +58,9 @@ def test_execute_workflow_add_content_and_publish( ], ) - finished = _wait_for_workflow(workflow_bindings.WorkflowsApi, workflow.pulp_href) + finished = monitor_workflow(workflow.pulp_href) # ---- Workflow-level assertions. - assert finished.state == "completed", ( - f"Workflow state={finished.state!r} error={finished.error!r}" - ) assert finished.error is None assert finished.started_at is not None assert finished.finished_at is not None diff --git a/pulp_workflow/tests/functional/api/test_workflow_callbacks.py b/pulp_workflow/tests/functional/api/test_workflow_callbacks.py new file mode 100644 index 0000000..ff16fc3 --- /dev/null +++ b/pulp_workflow/tests/functional/api/test_workflow_callbacks.py @@ -0,0 +1,207 @@ +"""End-to-end test that runs a Workflow with CallbackServices. + +The workflow syncs a file repository. It has callbacks attached for the ``completed`` and +``finished`` lifecycle events; both should fire and their dispatched tasks should complete +successfully. Workflow context (the workflow's name, state, and labels) is exposed to the callback +script via environment variables; the test sets a ``email=user@example.com`` label to exercise +the ``PULP_WORKFLOW_LABEL_EMAIL`` path. +""" + +import time +import uuid + +import pytest + +from pulpcore.plugin.util import extract_pk + +POLL_INTERVAL_SECONDS = 2.0 + + +def test_workflow_with_callback_on_sync( + workflow_bindings, + pulpcore_bindings, + file_bindings, + file_repo, + file_remote_factory, + basic_manifest_path, + workflow_factory, + callback_service_factory, + monitor_task, + monitor_workflow, +): + """A workflow that syncs a file repo and notifies a callback on completion. + + Mirrors the user-story in the issue: an admin registers a CallbackService pointing at a script + (here ``/bin/echo``) and attaches it to a workflow that syncs a file repo. The workflow's + email recipient lives in a ``email`` label and is exposed to the callback as + ``PULP_WORKFLOW_LABEL_EMAIL``. + """ + repo = file_repo + remote = file_remote_factory(manifest_path=basic_manifest_path, policy="immediate") + + callback_service = callback_service_factory(script="/bin/echo") + + workflow = workflow_factory( + pulp_labels={"email": "user@example.com"}, + tasks=[ + { + "task_name": "pulp_file.app.tasks.synchronizing.synchronize", + "task_kwargs": [ + { + "kwarg_key": "repository_pk", + "value": extract_pk(repo.pulp_href), + }, + { + "kwarg_key": "remote_pk", + "value": extract_pk(remote.pulp_href), + }, + {"kwarg_key": "mirror", "value": False}, + ], + "reserved_resources": [repo.pulp_href], + }, + ], + callbacks=[ + { + "callback_service": callback_service.pulp_href, + "callback_type": "completed", + }, + { + "callback_service": callback_service.pulp_href, + "callback_type": "finished", + }, + ], + ) + + finished = monitor_workflow(workflow.pulp_href) + + # ---- Workflow-level assertions. + assert finished.pulp_labels == {"email": "user@example.com"} + assert len(finished.tasks) == 1 + + # ---- The sync task itself ran. + sync_task = pulpcore_bindings.TasksApi.read(finished.tasks[0].dispatched_task) + assert sync_task.state == "completed" + assert sync_task.name == "pulp_file.app.tasks.synchronizing.synchronize" + + # ---- Both callbacks fired and completed. + assert len(finished.callbacks) == 2 + callback_types = sorted(cb.callback_type for cb in finished.callbacks) + assert callback_types == ["completed", "finished"] + for cb in finished.callbacks: + assert cb.callback_service == callback_service.pulp_href, ( + f"Unexpected callback_service: {cb.callback_service!r}" + ) + assert cb.dispatched_task is not None, f"Callback {cb.callback_type} was not dispatched" + # monitor_task raises PulpTaskError if the task ends in any non-completed state. + cb_task = monitor_task(cb.dispatched_task) + assert cb_task.name == "pulp_workflow.app.tasks.run_callback" + + +def test_workflow_callback_fires_on_cancel( + workflow_bindings, + pulpcore_bindings, + workflow_factory, + callback_service_factory, + monitor_task, +): + """A canceled-before-start workflow fires its CANCELED + FINISHED callbacks.""" + from datetime import datetime, timedelta, timezone + + callback_service = callback_service_factory(script="/bin/echo") + + # Schedule far enough in the future that we can cancel before it starts. + start_time = (datetime.now(timezone.utc) + timedelta(hours=1)).isoformat() + workflow = workflow_factory( + start_time=start_time, + callbacks=[ + { + "callback_service": callback_service.pulp_href, + "callback_type": "canceled", + }, + { + "callback_service": callback_service.pulp_href, + "callback_type": "finished", + }, + # 'completed' should NOT fire when the workflow is canceled. + { + "callback_service": callback_service.pulp_href, + "callback_type": "completed", + }, + ], + ) + + canceled = workflow_bindings.WorkflowsApi.workflows_cancel( + workflow.pulp_href, {"state": "canceled"} + ) + assert canceled.state == "canceled" + + # Re-fetch so we see the dispatched_task hrefs the cancel set asynchronously. + fired_types = set() + deadline = time.monotonic() + 30 + while time.monotonic() < deadline: + latest = workflow_bindings.WorkflowsApi.read(workflow.pulp_href) + fired_types = {cb.callback_type for cb in latest.callbacks if cb.dispatched_task} + if {"canceled", "finished"}.issubset(fired_types): + break + time.sleep(POLL_INTERVAL_SECONDS) + + assert {"canceled", "finished"}.issubset(fired_types), ( + f"Expected canceled+finished callbacks to fire, only fired: {fired_types}" + ) + # The 'completed' callback must not have fired. + completed_cb = next(cb for cb in latest.callbacks if cb.callback_type == "completed") + assert completed_cb.dispatched_task is None + + # The fired callbacks should run to completion. + for cb in latest.callbacks: + if not cb.dispatched_task: + continue + monitor_task(cb.dispatched_task) + + +def test_workflow_callback_duplicate_type_rejected(workflow_bindings, callback_service_factory): + """Two callbacks for the same (callback_service, callback_type) on one workflow are rejected.""" + callback_service = callback_service_factory() + with pytest.raises(workflow_bindings.ApiException) as exc: + workflow_bindings.WorkflowsApi.create( + { + "name": str(uuid.uuid4()), + "tasks": [{"task_name": "pulpcore.app.tasks.orphan_cleanup"}], + "callbacks": [ + { + "callback_service": callback_service.pulp_href, + "callback_type": "completed", + }, + { + "callback_service": callback_service.pulp_href, + "callback_type": "completed", + }, + ], + } + ) + assert exc.value.status == 400 + + +def test_workflow_callback_invalid_type_rejected(workflow_bindings, callback_service_factory): + """A callback_type outside of the documented choices is rejected. + + The bindings declare ``callback_type`` as a closed enum, so an invalid value is caught + client-side by pydantic before the request is sent. We accept either that pydantic + ``ValidationError`` or a server-side 400 (when called via raw HTTP). + """ + from pydantic import ValidationError + + callback_service = callback_service_factory() + with pytest.raises((ValidationError, workflow_bindings.ApiException)): + workflow_bindings.WorkflowsApi.create( + { + "name": str(uuid.uuid4()), + "tasks": [{"task_name": "pulpcore.app.tasks.orphan_cleanup"}], + "callbacks": [ + { + "callback_service": callback_service.pulp_href, + "callback_type": "not-a-real-state", + }, + ], + } + ) diff --git a/pulp_workflow/tests/unit/test_models.py b/pulp_workflow/tests/unit/test_models.py new file mode 100644 index 0000000..908139a --- /dev/null +++ b/pulp_workflow/tests/unit/test_models.py @@ -0,0 +1,76 @@ +from types import SimpleNamespace + +import pytest +from django.core.exceptions import ImproperlyConfigured +from django.test import override_settings + +from pulp_workflow.app.models import CallbackService + +# Constructing CallbackService triggers Django's pulp_domain default (get_domain_pk), +# which queries the DB. These tests don't read/write rows, but they need DB access. +pytestmark = pytest.mark.django_db + + +def _stub_workflow(): + return SimpleNamespace( + pk="00000000-0000-0000-0000-000000000001", + name="my-workflow", + state="completed", + pulp_labels={"email": "user@example.com"}, + ) + + +def test_callback_env_default_omits_pk_and_labels(): + """Default exposes only name+state; pk and labels are NOT leaked.""" + with override_settings(WORKFLOW_CALLBACK_FIELDS=["name", "state"]): + env = CallbackService(name="cb", script="/bin/echo")._env(_stub_workflow()) + assert env["PULP_WORKFLOW_NAME"] == "my-workflow" + assert env["PULP_WORKFLOW_STATE"] == "completed" + assert "PULP_WORKFLOW_PK" not in env + assert "PULP_WORKFLOW_LABELS" not in env + assert "PULP_WORKFLOW_LABEL_EMAIL" not in env + + +def test_callback_env_labels_field_exposes_per_label_vars(): + """Opting in to ``labels`` exposes both the JSON view and per-label vars.""" + with override_settings(WORKFLOW_CALLBACK_FIELDS=["labels"]): + env = CallbackService(name="cb", script="/bin/echo")._env(_stub_workflow()) + assert "user@example.com" in env["PULP_WORKFLOW_LABELS"] + assert env["PULP_WORKFLOW_LABEL_EMAIL"] == "user@example.com" + assert "PULP_WORKFLOW_NAME" not in env + + +def test_callback_env_unknown_field_raises(): + """A typo in the setting fails loudly rather than silently leaking/dropping data.""" + with override_settings(WORKFLOW_CALLBACK_FIELDS=["secrets"]): + with pytest.raises(ImproperlyConfigured, match="secrets"): + CallbackService(name="cb", script="/bin/echo")._env(_stub_workflow()) + + +def test_callback_env_labels_key_exposes_only_that_label(): + """``labels:`` exposes a single label and does NOT leak the JSON view or others.""" + workflow = SimpleNamespace( + pk="x", + name="w", + state="completed", + pulp_labels={"email": "user@example.com", "secret": "shh"}, + ) + with override_settings(WORKFLOW_CALLBACK_FIELDS=["labels:email"]): + env = CallbackService(name="cb", script="/bin/echo")._env(workflow) + assert env["PULP_WORKFLOW_LABEL_EMAIL"] == "user@example.com" + assert "PULP_WORKFLOW_LABEL_SECRET" not in env + assert "PULP_WORKFLOW_LABELS" not in env + + +def test_callback_env_labels_key_missing_label_is_empty_string(): + """Requested label that the workflow doesn't have still gets a (empty) env var.""" + with override_settings(WORKFLOW_CALLBACK_FIELDS=["labels:absent"]): + env = CallbackService(name="cb", script="/bin/echo")._env(_stub_workflow()) + assert env["PULP_WORKFLOW_LABEL_ABSENT"] == "" + + +def test_callback_env_labels_empty_key_raises(): + """``labels:`` with no key is a typo, not a valid request.""" + with override_settings(WORKFLOW_CALLBACK_FIELDS=["labels:"]): + with pytest.raises(ImproperlyConfigured, match="labels:"): + CallbackService(name="cb", script="/bin/echo")._env(_stub_workflow()) diff --git a/pulp_workflow/tests/unit/test_viewsets.py b/pulp_workflow/tests/unit/test_viewsets.py index 8a8794f..926c101 100644 --- a/pulp_workflow/tests/unit/test_viewsets.py +++ b/pulp_workflow/tests/unit/test_viewsets.py @@ -1,4 +1,4 @@ -from pulp_workflow.app.viewsets import WorkflowViewSet +from pulp_workflow.app.viewsets import CallbackServiceViewSet, WorkflowViewSet def test_access_policy_requires_view_for_read(): @@ -36,3 +36,35 @@ def test_locked_roles(): # has been removed. for perms in roles.values(): assert "workflow.delete_workflow" not in perms + + +def test_callback_service_access_policy_requires_view_for_read(): + """CallbackService read actions require view_callbackservice permission.""" + policy = CallbackServiceViewSet.DEFAULT_ACCESS_POLICY + read_stmt = policy["statements"][0] + assert set(read_stmt["action"]) == {"list", "retrieve", "my_permissions"} + assert "view_callbackservice" in read_stmt["condition"] + + +def test_callback_service_access_policy_distinguishes_write_actions(): + """add/change/delete on CallbackService each require their own permission.""" + statements = { + tuple(s["action"]): s for s in CallbackServiceViewSet.DEFAULT_ACCESS_POLICY["statements"] + } + assert "add_callbackservice" in statements[("create",)]["condition"] + assert "change_callbackservice" in statements[("update", "partial_update")]["condition"] + assert "delete_callbackservice" in statements[("destroy",)]["condition"] + + +def test_callback_service_locked_roles(): + """Admin and viewer roles are defined for CallbackService.""" + roles = CallbackServiceViewSet.LOCKED_ROLES + assert "workflow.callbackservice_admin" in roles + assert "workflow.callbackservice_viewer" in roles + admin_perms = roles["workflow.callbackservice_admin"] + assert "workflow.add_callbackservice" in admin_perms + assert "workflow.change_callbackservice" in admin_perms + assert "workflow.delete_callbackservice" in admin_perms + assert "workflow.view_callbackservice" in admin_perms + assert "workflow.manage_roles_callbackservice" in admin_perms + assert roles["workflow.callbackservice_viewer"] == ["workflow.view_callbackservice"] diff --git a/template_config.yml b/template_config.yml index 08eca54..5486af1 100644 --- a/template_config.yml +++ b/template_config.yml @@ -40,6 +40,10 @@ pulp_env_s3: {} pulp_scheme: "https" pulp_settings: api_root: "/pulp/" + WORKFLOW_CALLBACK_FIELDS: + - "name" + - "state" + - "labels:email" pulp_settings_azure: MEDIA_ROOT: "" STORAGES: