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: 1 addition & 1 deletion .github/workflows/scripts/before_install.sh
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
5 changes: 5 additions & 0 deletions CHANGES/10.feature
Original file line number Diff line number Diff line change
@@ -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.
Comment thread
daviddavis marked this conversation as resolved.
18 changes: 11 additions & 7 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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/<pk>/` | Retrieve a workflow |
| PATCH | `/pulp/api/v3/workflows/<pk>/` | 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/<pk>/` | Retrieve a workflow |
| PATCH | `/pulp/api/v3/workflow/workflows/<pk>/` | 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

Expand Down Expand Up @@ -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`).
state (`completed`, `failed`, or `canceled`).
132 changes: 132 additions & 0 deletions docs/demo/README.md
Original file line number Diff line number Diff line change
@@ -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
Comment thread
daviddavis marked this conversation as resolved.
```

## 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/ <<JSON | jq -r .pulp_href
{
"name": "demo-workflow",
"tasks": [
{
"task_name": "pulp_file.app.tasks.synchronizing.synchronize",
"reserved_resources": ["$REPO_HREF"],
"task_kwargs": [
{"kwarg_key": "repository_pk", "value": "$REPO_PK"},
{"kwarg_key": "remote_pk", "value": "$REMOTE_PK"},
{"kwarg_key": "mirror", "value": false}
]
},
{
"task_name": "pulp_file.app.tasks.publish",
"reserved_resources": ["$REPO_HREF"],
"task_kwargs": [
{"kwarg_key": "repository_version_pk",
"content_type": "core.repositoryversion"}
]
}
],
"callbacks": [
{"callback_service": "$CB_HREF", "callback_type": "finished"}
]
}
JSON
)
echo "WF_HREF=$WF_HREF"
```

## 5. Watch it run

```bash
while :; do
STATE=$(http :5001"$WF_HREF" | jq -r .state)
echo "state=$STATE"
case "$STATE" in
completed|failed|canceled) break ;;
esac
sleep 2
done

http :5001"$WF_HREF" | jq '{state, started_at, finished_at, error,
tasks: [.tasks[] | {task_name, dispatched_task}],
callbacks: [.callbacks[] | {callback_type, dispatched_task}]}'
```

On success: `state: "completed"`, every `dispatched_task` non-null, and a new
`RepositoryVersion` + `Publication`:

```bash
pulp file repository version list --repository "$REPO_NAME"
pulp file publication list --repository "$REPO_NAME"
```

The callback task ending in `completed` means `curl -fsS` got a 2xx from the
webhook:

```bash
CB_TASK_HREF=$(http :5001"$WF_HREF" | jq -r '.callbacks[0].dispatched_task')
http :5001"$CB_TASK_HREF" | jq '{name, state, error}'
```
13 changes: 13 additions & 0 deletions docs/demo/notify.sh
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
#!/bin/bash
# Demo callback script for `pulp_workflow`. Invoked by a `CallbackService` as a
# subprocess on the Pulp worker; reads workflow context from `PULP_WORKFLOW_*`
# env vars and `NOTIFY_WEBHOOK` from the worker's environment (set via
# oci_env/compose.env in the dev stack).
set -euo pipefail
: "${NOTIFY_WEBHOOK:?NOTIFY_WEBHOOK must be set in the worker environment}"
PAYLOAD=$(jq -nc \
--arg name "${PULP_WORKFLOW_NAME:-?}" \
--arg state "${PULP_WORKFLOW_STATE:-?}" \
--arg cid "${CORRELATION_ID:-?}" \
'{content: ("Workflow \($name) finished in state \($state).")}')
curl -fsS -H 'Content-Type: application/json' -d "$PAYLOAD" "$NOTIFY_WEBHOOK" >/dev/null
Original file line number Diff line number Diff line change
@@ -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),
),
]
Loading