JavaScript Worker

This guide shows you how to build a Priostack job worker in JavaScript/Node.js using the built-in fetch API (Node 18+).

Prerequisites

Activation answers at once (there is no long polling), so the worker pauses 2 seconds after an empty poll. fetch does not reject on a 4xx or 5xx status, so every call, including fail, checks resp.ok: activate answers 200, complete and fail answer 204. Keys are opaque strings. On SIGTERM the worker stops polling, waits for in-flight jobs with a bounded grace period, then hands unfinished jobs back with their retries unchanged.

Complete Worker Implementation

// worker.mjs - Priostack Job Worker (Node.js / ESM)

const BASE_URL = "https://priostack.com";
const WORKER_TYPE = "payment-processor";
const MAX_JOBS = 5;
const IDLE_WAIT_MS = 2_000;       // pause after an empty poll: activation answers at once
const SHUTDOWN_GRACE_MS = 20_000; // time in-flight jobs get to finish on shutdown

const apiKey = process.env.PRIOSTACK_API_KEY;
if (!apiKey) {
  console.error("PRIOSTACK_API_KEY environment variable is required");
  process.exit(1);
}

const headers = {
  "X-API-Key": apiKey,
  "Content-Type": "application/json",
};

const sleep = (ms) => new Promise((resolve) => setTimeout(resolve, ms));

// ── API helpers ────────────────────────────────────────────────

// fetch does not reject on 4xx/5xx, so every call checks resp.ok.
// Activate answers 200; complete and fail answer 204 with an empty body.
async function post(path, body) {
  const resp = await fetch(`${BASE_URL}${path}`, {
    method: "POST",
    headers,
    body: JSON.stringify(body),
    signal: AbortSignal.timeout(15_000),
  });
  const text = await resp.text();
  if (!resp.ok) {
    const err = new Error(`${path} failed ${resp.status}: ${text}`);
    err.status = resp.status;
    throw err;
  }
  return text ? JSON.parse(text) : null;
}

async function activateJobs() {
  const data = await post("/api/v1/jobs/activate", {
    type: WORKER_TYPE,
    maxJobsToActivate: MAX_JOBS,
  });
  return data?.jobs ?? [];
}

async function completeJob(jobKey, variables) {
  await post(`/api/v1/jobs/${encodeURIComponent(jobKey)}/complete`, { variables });
}

// 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.
async function failJob(jobKey, retries, errorMessage) {
  await post(`/api/v1/jobs/${encodeURIComponent(jobKey)}/fail`, { retries, errorMessage });
}

// ── Business logic ─────────────────────────────────────────────

async function processPayment(job) {
  const { orderId, amount, currency = "EUR" } = job.variables ?? {};
  console.log(`Processing payment for order ${orderId}, amount ${amount} ${currency}`);

  // TODO: Call your actual payment gateway here
  // const intent = await stripe.paymentIntents.create({ amount: Math.round(amount * 100), currency });

  // Simulate async processing
  await sleep(100);

  return {
    transactionId: `txn_${orderId}_${Date.now()}`,
    paymentStatus: "success",
    processedAt: new Date().toISOString(),
  };
}

// ── Worker loop ────────────────────────────────────────────────

// report runs one complete or fail call. A 404 means the job is gone:
// unknown, already completed, or its instance was cancelled.
async function report(call, jobKey) {
  try {
    await call();
    return true;
  } catch (err) {
    if (err.status === 404) {
      console.warn(`Job ${jobKey} not found; it was probably already applied`);
    } else {
      console.error(`Reporting job ${jobKey} failed:`, err.message);
    }
    return false;
  }
}

async function handleJob(job) {
  const { key: jobKey, retries } = job; // key is an opaque string
  let result;
  try {
    result = await processPayment(job);
  } catch (err) {
    console.error(`Job ${jobKey} failed:`, err.message);
    // Send one less than the job's retries to use up one attempt.
    await report(() => failJob(jobKey, retries - 1, err.message), jobKey);
    return;
  }
  if (await report(() => completeJob(jobKey, result), jobKey)) {
    console.log(`Job ${jobKey} completed successfully`);
  }
}

const inFlight = new Map(); // job key -> job, while its handler runs
let running = true;

process.on("SIGTERM", () => { console.log("SIGTERM received. Shutting down..."); running = false; });
process.on("SIGINT",  () => { console.log("SIGINT received. Shutting down...");  running = false; });

async function main() {
  console.log(`Worker started. Polling for jobs of type "${WORKER_TYPE}"...`);

  while (running) {
    let jobs = [];
    try {
      jobs = await activateJobs();
    } catch (err) {
      console.error("Activate failed:", err.message);
      if (err.status === 429) {
        await sleep(60_000); // the limit is 300 requests per 60 s per IP
        continue;
      }
    }
    for (const job of jobs) {
      console.log(`Activated job ${job.key} (instance ${job.processInstanceKey})`);
      // Process jobs concurrently, and remember them until they finish.
      inFlight.set(job.key, job);
      handleJob(job)
        .catch((err) => console.error("Unhandled job error:", err))
        .finally(() => inFlight.delete(job.key));
    }
    // No long polling: an empty answer comes back at once, so pause.
    if (jobs.length === 0 && running) await sleep(IDLE_WAIT_MS);
  }

  // Stop polling, then give in-flight jobs a bounded time to finish.
  const deadline = Date.now() + SHUTDOWN_GRACE_MS;
  while (inFlight.size > 0 && Date.now() < deadline) await sleep(200);

  // Hand back what is still running, without using up an attempt. The server
  // never expires a lock, so a job left active would stay stranded.
  for (const job of inFlight.values()) {
    await report(() => failJob(job.key, job.retries, "worker shutting down"), job.key);
  }
  console.log("Worker stopped.");
  process.exit(0);
}

main().catch((err) => {
  console.error("Fatal:", err);
  process.exit(1);
});

Running the Worker

# Set your API key
export PRIOSTACK_API_KEY=ps_your_key_here

# Run with Node.js (ESM)
node worker.mjs

# Or with tsx (TypeScript)
npx tsx worker.ts

# Output:
# Worker started. Polling for jobs of type "payment-processor"...
# Activated job a99f579d62ae71f8764f (instance run_1ce49bbc5260b39423bd)
# Processing payment for order ORD-001, amount 149.99 EUR
# Job a99f579d62ae71f8764f completed successfully

TypeScript Version

To use TypeScript, add type annotations:

interface Job {
  key: string;                  // opaque, 20 hex characters
  type: string;
  processInstanceKey: string;   // opaque, currently "run_" + 20 hex
  processDefinitionKey: string; // empty for CMMN case jobs
  variables: Record<string, unknown>;
  retries: number;
  deadline: number;             // Unix ms, informational only
  createdAt: string;
}

interface ActivateResponse {
  jobs: Job[];
}

async function activateJobs(): Promise<Job[]> {
  // same implementation, typed
}

Dockerfile

FROM node:20-alpine
WORKDIR /app
COPY worker.mjs .
CMD ["node", "worker.mjs"]

Production Tips

ConcernRecommendation
ConcurrencyJobs from one activation already run concurrently. For CPU-heavy work, use a worker pool (e.g., piscina) and keep maxJobsToActivate at what the pool can take.
Graceful shutdownTrack in-flight jobs, wait for them with a bounded grace period on SIGTERM, then fail unfinished ones with retries unchanged, as the code above does. The server never expires a lock.
Error classificationDistinguish network errors (retry) from business errors (complete the job with an outcome variable and branch on it: there is no REST call to throw a BPMN error) from programming errors (log and alert).
MonitoringUse --inspect for debugging. Export metrics to Prometheus using prom-client.