Go Worker
This guide shows you how to build a production-ready Priostack job worker in Go using only the standard library - no SDK required.
Prerequisites
- Go 1.21 or later
- A Priostack API key (set as
PRIOSTACK_API_KEYenvironment variable) - A deployed BPMN process with a Service Task of type
payment-processor
The worker polls POST /api/v1/jobs/activate, which answers at once (there is no long polling), so it pauses 2 seconds after an empty or failed poll. It treats any 2xx as success (activate answers 200, complete and fail answer 204), decodes every key as a string, and on shutdown stops polling, waits for in-flight jobs with a bounded grace period, then hands unfinished jobs back with their retries unchanged.
Complete Worker Implementation
package main
import (
"bytes"
"context"
"encoding/json"
"errors"
"fmt"
"io"
"log"
"net/http"
"net/url"
"os"
"os/signal"
"sync"
"syscall"
"time"
)
const (
baseURL = "https://priostack.com"
workerType = "payment-processor"
maxJobs = 5
idleWait = 2 * time.Second // pause after an empty or failed poll
shutdownGrace = 20 * time.Second // time in-flight jobs get to finish
)
// Job is one entry of the activate response. Every key is an opaque string:
// declaring them as int64 makes the whole response fail to decode.
type Job struct {
Key string `json:"key"`
Type string `json:"type"`
ProcessInstanceKey string `json:"processInstanceKey"`
ProcessDefinitionKey string `json:"processDefinitionKey"`
Variables map[string]interface{} `json:"variables"`
Retries int `json:"retries"`
Deadline int64 `json:"deadline"` // Unix ms, informational only
}
type ActivateResponse struct {
Jobs []Job `json:"jobs"`
}
// StatusError is an answer outside 2xx. A 404 on complete or fail means the
// job is gone: unknown, already completed, or its instance was cancelled.
type StatusError struct {
Path string
Code int
Body string
}
func (e *StatusError) Error() string {
return fmt.Sprintf("%s: status %d: %s", e.Path, e.Code, e.Body)
}
type PriostackClient struct {
httpClient *http.Client
apiKey string
baseURL string
}
func NewClient(apiKey, baseURL string) *PriostackClient {
return &PriostackClient{
httpClient: &http.Client{Timeout: 15 * time.Second},
apiKey: apiKey,
baseURL: baseURL,
}
}
// post sends one JSON call and returns the response body. Activate answers
// 200 and complete and fail answer 204, so any 2xx is success.
func (c *PriostackClient) post(ctx context.Context, path string, body interface{}) ([]byte, error) {
b, err := json.Marshal(body)
if err != nil {
return nil, fmt.Errorf("marshal: %w", err)
}
req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.baseURL+path, bytes.NewReader(b))
if err != nil {
return nil, fmt.Errorf("new request: %w", err)
}
req.Header.Set("X-API-Key", c.apiKey)
req.Header.Set("Content-Type", "application/json")
resp, err := c.httpClient.Do(req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
data, err := io.ReadAll(io.LimitReader(resp.Body, 10<<20))
if err != nil {
return nil, fmt.Errorf("read body: %w", err)
}
if resp.StatusCode/100 != 2 {
return nil, &StatusError{Path: path, Code: resp.StatusCode, Body: string(data)}
}
return data, nil
}
func (c *PriostackClient) ActivateJobs(ctx context.Context) ([]Job, error) {
data, err := c.post(ctx, "/api/v1/jobs/activate", map[string]interface{}{
"type": workerType,
"maxJobsToActivate": maxJobs,
})
if err != nil {
return nil, fmt.Errorf("activate: %w", err)
}
var result ActivateResponse
if err := json.Unmarshal(data, &result); err != nil {
return nil, fmt.Errorf("decode: %w", err)
}
return result.Jobs, nil
}
func (c *PriostackClient) CompleteJob(ctx context.Context, jobKey string, variables map[string]interface{}) error {
path := "/api/v1/jobs/" + url.PathEscape(jobKey) + "/complete"
_, err := c.post(ctx, path, map[string]interface{}{"variables": variables})
return err
}
// FailJob stores retries as the job's attempts left: the server does not
// decrement it. Above 0 the job is re-queued at once; 0 raises an incident.
func (c *PriostackClient) FailJob(ctx context.Context, jobKey string, retries int, errMsg string) error {
path := "/api/v1/jobs/" + url.PathEscape(jobKey) + "/fail"
_, err := c.post(ctx, path, map[string]interface{}{
"retries": retries,
"errorMessage": errMsg,
})
return err
}
// processPayment is your actual business logic. It returns early when ctx is
// cancelled, so a shutdown can interrupt it.
func processPayment(ctx context.Context, job Job) (map[string]interface{}, error) {
orderID, _ := job.Variables["orderId"].(string)
amount, _ := job.Variables["amount"].(float64)
log.Printf("Processing payment for order %s, amount %.2f", orderID, amount)
// TODO: call your actual payment gateway here, passing ctx.
select {
case <-time.After(100 * time.Millisecond):
case <-ctx.Done():
return nil, ctx.Err()
}
return map[string]interface{}{
"transactionId": "txn_" + orderID,
"paymentStatus": "success",
"processedAt": time.Now().UTC().Format(time.RFC3339),
}, nil
}
// handle runs one job and reports the outcome. The report calls get their own
// short context, so a job that finishes during shutdown is still reported.
func handle(work context.Context, client *PriostackClient, job Job) {
result, err := processPayment(work, job)
report, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
switch {
case err == nil:
if err := client.CompleteJob(report, job.Key, result); err != nil {
var se *StatusError
if errors.As(err, &se) && se.Code == http.StatusNotFound {
// Gone: probably applied by an earlier call whose answer was lost.
log.Printf("Job %s not found; check instance %s", job.Key, job.ProcessInstanceKey)
return
}
log.Printf("Failed to complete job %s: %v", job.Key, err)
return
}
log.Printf("Job %s completed successfully", job.Key)
case work.Err() != nil:
// Shutdown interrupted the work: hand the job back without using up
// an attempt, so it can be activated again at once.
if err := client.FailJob(report, job.Key, job.Retries, "worker shutting down"); err != nil {
log.Printf("Failed to release job %s: %v", job.Key, err)
return
}
log.Printf("Job %s released", job.Key)
default:
log.Printf("Job %s failed: %v", job.Key, err)
if err := client.FailJob(report, job.Key, job.Retries-1, err.Error()); err != nil {
log.Printf("Failed to report job failure: %v", err)
}
}
}
func runWorker(ctx context.Context, client *PriostackClient) {
// work outlives ctx by the grace period, so in-flight jobs can finish.
work, stopWork := context.WithCancel(context.Background())
defer stopWork()
var inFlight sync.WaitGroup
log.Printf("Worker started. Polling for jobs of type %q...", workerType)
for ctx.Err() == nil {
// The poll uses work, not ctx: abandoning a response the server has
// already sent would strand the jobs it activated.
jobs, err := client.ActivateJobs(work)
pause := idleWait
if err != nil {
log.Printf("Error activating jobs: %v", err)
var se *StatusError
if errors.As(err, &se) && se.Code == http.StatusTooManyRequests {
pause = time.Minute // the limit is 300 requests per 60 s per IP
}
}
for _, job := range jobs {
inFlight.Add(1)
go func(job Job) {
defer inFlight.Done()
log.Printf("Activated job %s (instance %s)", job.Key, job.ProcessInstanceKey)
handle(work, client, job)
}(job)
}
if len(jobs) == 0 {
// No long polling: an empty answer comes back at once, so wait.
select {
case <-time.After(pause):
case <-ctx.Done():
}
}
}
// Stop activating, then give in-flight jobs a bounded time to finish.
log.Println("Worker shutting down. Waiting for in-flight jobs...")
done := make(chan struct{})
go func() {
inFlight.Wait()
close(done)
}()
select {
case <-done:
case <-time.After(shutdownGrace):
log.Println("Grace period over. Releasing unfinished jobs...")
stopWork()
<-done
}
}
func main() {
apiKey := os.Getenv("PRIOSTACK_API_KEY")
if apiKey == "" {
log.Fatal("PRIOSTACK_API_KEY environment variable is required")
}
client := NewClient(apiKey, baseURL)
ctx, cancel := signal.NotifyContext(context.Background(),
os.Interrupt, syscall.SIGTERM)
defer cancel()
runWorker(ctx, client)
log.Println("Worker stopped.")
}
Running the Worker
# Set your API key
export PRIOSTACK_API_KEY=ps_your_key_here
# Run the worker
go run worker.go
# Output:
# 2026/05/01 14:00:00 Worker started. Polling for jobs of type "payment-processor"...
# 2026/05/01 14:00:05 Activated job a99f579d62ae71f8764f (instance run_1ce49bbc5260b39423bd)
# 2026/05/01 14:00:05 Processing payment for order ORD-001, amount 149.99
# 2026/05/01 14:00:05 Job a99f579d62ae71f8764f completed successfully
Production Deployment Tips
| Concern | Recommendation |
|---|---|
| Scaling | Run multiple instances of the worker binary. Each instance polls independently. Use Kubernetes HPA to scale based on pending job count. |
| Observability | Use structured logging (log/slog in Go 1.21+). Emit metrics for job processing time, error rate, and throughput. |
| Retries | Send job.Retries - 1 to use up one attempt: the server stores the number you send and does not decrement it. A job failed with retries above 0 can be activated again immediately; there is no retry delay, so wait in the worker first if the downstream system needs time. |
| Timeouts | The job deadline (Unix milliseconds, activation + 5 minutes) is informational: nothing reclaims the job when it passes. Bound your own processing time with a context deadline, and model a timer boundary event on the task if the process must move on. |
| Shutdown | The server never expires a lock, so a worker that exits holding jobs strands them until a server restart. Wait for in-flight jobs and release unfinished ones with their retries unchanged, as runWorker does. |
| Secrets | Inject PRIOSTACK_API_KEY via Kubernetes Secrets or a secrets manager. Never bake it into the Docker image. |
Dockerfile
FROM golang:1.21-alpine AS builder
WORKDIR /app
COPY . .
RUN go build -o worker ./worker.go
FROM alpine:3.19
COPY --from=builder /app/worker /worker
ENTRYPOINT ["/worker"]