Python Worker
This guide shows you how to build a Priostack job worker in Python using the requests library.
Prerequisites
- Python 3.10 or later
- The
requestslibrary:pip install requests - A Priostack API key set as
PRIOSTACK_API_KEYenvironment variable
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
| Concern | Recommendation |
|---|---|
| Concurrency | Run 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 handling | Distinguish 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 check | Run a simple HTTP server in a background thread to expose /health for Kubernetes liveness probes. |
| Secrets | Use environment variables or a secrets manager. Never hardcode credentials. |