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

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

ConcernRecommendation
ScalingRun multiple instances of the worker binary. Each instance polls independently. Use Kubernetes HPA to scale based on pending job count.
ObservabilityUse structured logging (log/slog in Go 1.21+). Emit metrics for job processing time, error rate, and throughput.
RetriesSend 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.
TimeoutsThe 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.
ShutdownThe 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.
SecretsInject 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"]