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
- Node.js 18 or later (built-in
fetchAPI) - 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. 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
| Concern | Recommendation |
|---|---|
| Concurrency | Jobs 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 shutdown | Track 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 classification | Distinguish 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). |
| Monitoring | Use --inspect for debugging. Export metrics to Prometheus using prom-client. |