Python Worker

This guide shows you how to build a Priostack job worker in Python using the requests library.

Prerequisites

Activation answers at once (there is no long polling), so the worker pauses 2 seconds after an empty poll. Every call checks its status with raise_for_status(): activate answers 200, complete and fail answer 204. Keys are opaque strings. Jobs are handled one at a time, so a shutdown signal lets the current job finish and hands any job still waiting in the batch back with its retries unchanged.

Complete Worker Implementation

#!/usr/bin/env python3
"""
Priostack Job Worker - Python implementation
"""
import logging
import os
import signal
import time
from datetime import datetime, timezone
from typing import Any, Callable

import requests

logging.basicConfig(
    level=logging.INFO,
    format="%(asctime)s %(levelname)s %(message)s",
)
log = logging.getLogger(__name__)

BASE_URL = "https://priostack.com"
WORKER_TYPE = "payment-processor"
MAX_JOBS = 5
IDLE_WAIT_S = 2  # pause after an empty poll: activation answers at once


class PriostackClient:
    """Every call raises on a status outside 2xx (raise_for_status).
    Activate answers 200; complete and fail answer 204 with an empty body."""

    def __init__(self, api_key: str, base_url: str = BASE_URL):
        self.session = requests.Session()
        self.session.headers.update({
            "X-API-Key": api_key,
            "Content-Type": "application/json",
        })
        self.base_url = base_url

    def activate_jobs(self) -> list[dict]:
        resp = self.session.post(
            f"{self.base_url}/api/v1/jobs/activate",
            json={"type": WORKER_TYPE, "maxJobsToActivate": MAX_JOBS},
            timeout=15,
        )
        resp.raise_for_status()
        return resp.json().get("jobs", [])

    def complete_job(self, job_key: str, variables: dict[str, Any]) -> None:
        resp = self.session.post(
            f"{self.base_url}/api/v1/jobs/{job_key}/complete",
            json={"variables": variables},
            timeout=10,
        )
        resp.raise_for_status()

    def fail_job(self, job_key: str, retries: int, error_message: str) -> None:
        # retries is stored as the attempts left; the server does not
        # decrement it. Above 0 the job is re-queued at once, 0 raises an incident.
        resp = self.session.post(
            f"{self.base_url}/api/v1/jobs/{job_key}/fail",
            json={"retries": retries, "errorMessage": error_message},
            timeout=10,
        )
        resp.raise_for_status()


def process_payment(job: dict) -> dict[str, Any]:
    """Your actual business logic goes here."""
    variables = job.get("variables", {})
    order_id = variables.get("orderId", "unknown")
    amount = variables.get("amount", 0)

    log.info("Processing payment for order %s, amount %.2f", order_id, amount)

    # TODO: Call your actual payment gateway
    # Example: response = stripe.PaymentIntent.create(amount=int(amount*100), currency="eur")

    # Simulate processing time
    time.sleep(0.1)

    return {
        "transactionId": f"txn_{order_id}_{int(time.time())}",
        "paymentStatus": "success",
        "processedAt": datetime.now(timezone.utc).isoformat(),
    }


class Worker:
    def __init__(self, client: PriostackClient):
        self.client = client
        self.running = True
        signal.signal(signal.SIGTERM, self._shutdown)
        signal.signal(signal.SIGINT, self._shutdown)

    def _shutdown(self, signum, frame):
        log.info("Shutdown signal received. Finishing the current job...")
        self.running = False

    def report(self, call: Callable[..., None], job_key: str, *args) -> bool:
        try:
            call(job_key, *args)
            return True
        except requests.exceptions.HTTPError as exc:
            if exc.response is not None and exc.response.status_code == 404:
                # The job is gone: unknown, already completed, or its instance
                # was cancelled. Read the instance to reconcile.
                log.warning("Job %s not found; it was probably already applied", job_key)
            else:
                log.error("Reporting job %s failed: %s", job_key, exc)
        except requests.exceptions.RequestException as exc:
            log.error("Reporting job %s failed: %s", job_key, exc)
        return False

    def handle_job(self, job: dict) -> None:
        job_key = job["key"]  # an opaque string, never a number
        try:
            result = process_payment(job)
        except Exception as exc:
            log.error("Job %s failed: %s", job_key, exc)
            # Send one less than the job's retries to use up one attempt.
            self.report(self.client.fail_job, job_key, job["retries"] - 1, str(exc))
            return
        if self.report(self.client.complete_job, job_key, result):
            log.info("Job %s completed successfully", job_key)

    def sleep(self, seconds: float) -> None:
        # Sleep in short steps so a shutdown signal is noticed quickly.
        end = time.monotonic() + seconds
        while self.running and time.monotonic() < end:
            time.sleep(max(0.0, min(0.5, end - time.monotonic())))

    def run(self) -> None:
        log.info("Worker started. Polling for jobs of type %r...", WORKER_TYPE)

        while self.running:
            try:
                jobs = self.client.activate_jobs()
            except requests.exceptions.HTTPError as exc:
                if exc.response is not None and exc.response.status_code == 429:
                    log.warning("Rate limited. Waiting 60s...")
                    self.sleep(60)
                else:
                    log.error("HTTP error: %s. Retrying in 5s...", exc)
                    self.sleep(5)
                continue
            except requests.exceptions.RequestException as exc:
                log.error("Network error: %s. Retrying in 5s...", exc)
                self.sleep(5)
                continue

            for job in jobs:
                if not self.running:
                    # Shutting down: hand the job back without using up an
                    # attempt. The server never expires a lock on its own.
                    self.report(self.client.fail_job, job["key"], job["retries"], "worker shutting down")
                    continue
                log.info("Activated job %s (instance %s)", job["key"], job.get("processInstanceKey"))
                self.handle_job(job)

            if not jobs:
                # No long polling: an empty answer comes back at once.
                self.sleep(IDLE_WAIT_S)

        log.info("Worker stopped.")


def main():
    api_key = os.environ.get("PRIOSTACK_API_KEY")
    if not api_key:
        raise SystemExit("PRIOSTACK_API_KEY environment variable is required")

    client = PriostackClient(api_key)
    Worker(client).run()


if __name__ == "__main__":
    main()

Running the Worker

# Install dependencies
pip install requests

# Set your API key
export PRIOSTACK_API_KEY=ps_your_key_here

# Run the worker
python worker.py

# Output:
# 2026-05-01 14:00:00 INFO Worker started. Polling for jobs of type 'payment-processor'...
# 2026-05-01 14:00:05 INFO Activated job a99f579d62ae71f8764f (instance run_1ce49bbc5260b39423bd)
# 2026-05-01 14:00:05 INFO Processing payment for order ORD-001, amount 149.99
# 2026-05-01 14:00:05 INFO Job a99f579d62ae71f8764f completed successfully

Dockerfile

FROM python:3.12-slim
WORKDIR /app
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
COPY worker.py .
CMD ["python", "worker.py"]

requirements.txt

requests==2.31.0

Production Tips

ConcernRecommendation
ConcurrencyRun multiple worker processes (use gunicorn or supervisor) or use threading.Thread for parallel job handling. With threads, join them on shutdown and fail unfinished jobs with retries unchanged: the server never expires a lock.
Error handlingDistinguish transient errors (network, timeout) from permanent errors. A business outcome cannot be thrown as a BPMN error over REST (there is no /error action): complete the job with an outcome variable and branch on it with a gateway.
Health checkRun a simple HTTP server in a background thread to expose /health for Kubernetes liveness probes.
SecretsUse environment variables or a secrets manager. Never hardcode credentials.