From AI prototype to credible information system.Explore Go Live Fintech
Engineering

Implementing the BPMN Service Task Worker Pattern in Go

When you model a BPMN 2.0 process that calls an external service - sending an email, charging a credit card, calling a third-party API - you reach for a Service Task. The execution engine needs something to actually run that task, and that something is a job worker. Getting the worker pattern right in Go (Golang) means your workflow engine stays responsive, your jobs are idempotent, and errors feed back into the process, where the model decides what happens next.

This tutorial walks through the complete BPMN service task worker pattern in Go using Priostack's REST API. By the end you will have a production-ready worker that polls for jobs, handles variables, retries on failure, reports business outcomes the process can route on, and shuts down without stranding a job.

What is a Service Task Worker?

In BPMN 2.0, a Service Task represents work performed by an automated system rather than a human. When the workflow engine reaches a service task, it creates a job and waits. A worker is a separate process that:

  1. Polls the engine for available jobs of a specific type (e.g., send-email).
  2. Holds the job: an activated job goes to this worker alone and stays active until it is completed or failed.
  3. Executes the business logic (call the email API).
  4. Reports success (complete) or failure (fail) back to the engine.

This asynchronous, polling-based architecture means your workers can be deployed independently, scaled horizontally, and replaced without redeploying the workflow engine. It is the same model Zeebe (Camunda 8) uses. Priostack follows the shape of Zeebe's job calls over HTTPS and JSON (no gRPC), so a Go worker is plain net/http code.

Setting Up the Priostack Worker SDK

Priostack exposes a REST job API. You do not need a gRPC dependency or a heavy SDK. A standard net/http client is sufficient. First, sign up at priostack.com/quickstart and get your API key.

Create a new Go module for your worker:

# Create worker project mkdir email-worker && cd email-worker go mod init example.com/email-worker

Define a minimal client that wraps the Priostack base URL and API key. Its post helper checks the status of every call and builds a fresh request each time, so a retry never resends a consumed body:

package main import ( "bytes" "context" "encoding/json" "fmt" "io" "log" "net/http" "net/url" "sync" "time" ) type Client struct { BaseURL string APIKey string HTTP *http.Client } func NewClient(baseURL, apiKey string) *Client { return &Client{ BaseURL: baseURL, APIKey: apiKey, HTTP: &http.Client{Timeout: 30 * time.Second}, } } // post sends one JSON request and checks the status. It builds a fresh // request (and body reader) on every call, so a caller can safely retry. func (c *Client) post(ctx context.Context, path string, payload any, out any) error { body, err := json.Marshal(payload) if err != nil { return err } req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.BaseURL+path, bytes.NewReader(body)) if err != nil { return err } req.Header.Set("Authorization", "Bearer "+c.APIKey) req.Header.Set("Content-Type", "application/json") resp, err := c.HTTP.Do(req) if err != nil { return err } defer resp.Body.Close() if resp.StatusCode/100 != 2 { // activate answers 200, complete and fail 204 msg, _ := io.ReadAll(io.LimitReader(resp.Body, 512)) return fmt.Errorf("POST %s: %d %s", path, resp.StatusCode, bytes.TrimSpace(msg)) } if out != nil { return json.NewDecoder(resp.Body).Decode(out) } return nil }

Now implement the ActivateJobs call. It answers at once with the jobs that are waiting, or an empty list; it never holds the request open, so the polling loop below pauses between empty polls.

// Every key is an opaque string ("a99f579d62ae71f8764f", "run_1ce4..."), never a number. type Job struct { Key string `json:"key"` Type string `json:"type"` ProcessInstanceKey string `json:"processInstanceKey"` Variables map[string]any `json:"variables"` Retries int `json:"retries"` Deadline int64 `json:"deadline"` // Unix ms, informational only } // ActivateJobs answers at once, with an empty list when nothing waits. // Always send the type, and an explicit maxJobsToActivate: there is no // server default, and an empty type activates jobs of every type. func (c *Client) ActivateJobs(ctx context.Context, jobType string, maxJobs int) ([]Job, error) { var result struct { Jobs []Job `json:"jobs"` } err := c.post(ctx, "/api/v1/jobs/activate", map[string]any{ "type": jobType, "maxJobsToActivate": maxJobs, }, &result) return result.Jobs, err }

Handling Variables and Output Mapping

Service tasks exchange data with the process via variables. Input variables are available on job.Variables; output variables are passed when completing the job.

// CompleteJob answers 204. A 404 means the job is gone (already completed, // failed or cancelled): read the instance to see where it stands. func (c *Client) CompleteJob(ctx context.Context, jobKey string, outputVars map[string]any) error { return c.post(ctx, "/api/v1/jobs/"+url.PathEscape(jobKey)+"/complete", map[string]any{"variables": outputVars}, nil) }

To read input variables safely, use type assertions with defaults:

// Reading typed variables from a job toEmail, _ := job.Variables["recipientEmail"].(string) subject, _ := job.Variables["emailSubject"].(string) if toEmail == "" { // variable missing - fail the job client.FailJob(ctx, job.Key, "recipientEmail variable is required", 0) return }

Error Handling and Retry Logic

Workers distinguish between two types of failure:

  • Retryable failure - transient error (network timeout, rate limit). Call FailJob with job.Retries - 1: the server stores the number you send and does not decrement it. A job failed with retries left is offered again at once; at zero the engine creates an incident.
  • Business outcome - a domain result such as a bounced email. There is no REST call to throw a BPMN error (/error answers 404), so complete the job with an output variable (e.g., emailBounced) and route on it with an exclusive gateway after the task.
// FailJob stores retries as the job's remaining attempts; the server does not // decrement it. retries > 0 puts the job back in the queue at once, and // retries <= 0 raises an incident. func (c *Client) FailJob(ctx context.Context, jobKey, reason string, retries int) error { return c.post(ctx, "/api/v1/jobs/"+url.PathEscape(jobKey)+"/fail", map[string]any{ "retries": retries, "errorMessage": reason, }, nil) }

Poll in a loop: pause 2 seconds after an empty poll, back off exponentially after an error, and run each job in its own goroutine. On shutdown, stop activating, wait for the jobs in flight, and release any job still held: the server never takes a job back by itself, so a worker that exits holding jobs strands them.

type JobHandler func(ctx context.Context, job Job) // Worker polls one job type and runs each job in its own goroutine. type Worker struct { Client *Client JobType string Handler JobHandler Drain time.Duration // how long shutdown waits for jobs in flight mu sync.Mutex inflight map[string]Job wg sync.WaitGroup } func (w *Worker) Run(ctx context.Context) { w.inflight = map[string]Job{} // Handlers get their own context so a shutdown signal does not cut off // the complete call of a job that is about to finish. work, stopWork := context.WithCancel(context.Background()) defer stopWork() backoff := time.Second for ctx.Err() == nil { jobs, err := w.Client.ActivateJobs(ctx, w.JobType, 10) if err != nil { if ctx.Err() == nil { log.Printf("activate: %v - retrying in %s", err, backoff) } sleep(ctx, backoff) backoff = min(backoff*2, 60*time.Second) continue } backoff = time.Second if len(jobs) == 0 { // Activation never waits: pause between empty polls, or a tight // loop runs into the 300 requests per minute limit. sleep(ctx, 2*time.Second) continue } for _, job := range jobs { w.mu.Lock() w.inflight[job.Key] = job w.mu.Unlock() w.wg.Add(1) go func(job Job) { defer w.wg.Done() w.Handler(work, job) w.mu.Lock() delete(w.inflight, job.Key) w.mu.Unlock() }(job) } } // Shutdown: no more activations. Wait for the jobs in flight, but not forever. done := make(chan struct{}) go func() { w.wg.Wait(); close(done) }() select { case <-done: return case <-time.After(w.Drain): } stopWork() // The server never takes a job back by itself, so release every job still // held. Unchanged retries re-queue it without consuming an attempt. w.mu.Lock() left := make([]Job, 0, len(w.inflight)) for _, job := range w.inflight { left = append(left, job) } w.mu.Unlock() for _, job := range left { rctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) if err := w.Client.FailJob(rctx, job.Key, "worker shutting down", job.Retries); err != nil { log.Printf("release job %s: %v", job.Key, err) } cancel() } } func sleep(ctx context.Context, d time.Duration) { select { case <-ctx.Done(): case <-time.After(d): } }

Complete Example: Email Notification Service Task

Here is a full worker that handles an email-notification service task. It reads recipient, subject, and body from process variables, calls a stand-in email API, and completes the job with messageSid and emailBounced output variables. With the code above in client.go and this in main.go, go vet passes.

package main import ( "context" "fmt" "log" "os" "os/signal" "syscall" "time" ) func handleEmailNotification(client *Client) JobHandler { return func(ctx context.Context, job Job) { to, _ := job.Variables["recipientEmail"].(string) subject, _ := job.Variables["emailSubject"].(string) body, _ := job.Variables["emailBody"].(string) if to == "" { // Retrying cannot help: raise an incident straight away. if err := client.FailJob(ctx, job.Key, "recipientEmail is required", 0); err != nil { log.Printf("fail job %s: %v", job.Key, err) } return } sid, bounced, err := sendEmail(ctx, to, subject, body) if err != nil { // Transient: consume one attempt. At zero the engine raises an incident. if ferr := client.FailJob(ctx, job.Key, fmt.Sprintf("send failed: %v", err), job.Retries-1); ferr != nil { log.Printf("fail job %s: %v", job.Key, ferr) } return } // Success, or a business outcome the process routes on with a gateway. if err := client.CompleteJob(ctx, job.Key, map[string]any{ "messageSid": sid, "emailBounced": bounced, "emailSentAt": time.Now().UTC().Format(time.RFC3339), }); err != nil { log.Printf("complete job %s: %v", job.Key, err) return } log.Printf("job %s completed: email to %s (sid=%s, bounced=%t)", job.Key, to, sid, bounced) } } // sendEmail stands in for your email provider's client. func sendEmail(ctx context.Context, to, subject, body string) (sid string, bounced bool, err error) { return fmt.Sprintf("msg-%d", time.Now().UnixNano()), false, ctx.Err() } func main() { client := NewClient( os.Getenv("PRIOSTACK_URL"), os.Getenv("PRIOSTACK_API_KEY"), ) ctx, cancel := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM) defer cancel() w := &Worker{ Client: client, JobType: "email-notification", Handler: handleEmailNotification(client), Drain: 30 * time.Second, } log.Println("Email worker started - polling for email-notification jobs") w.Run(ctx) log.Println("Worker stopped") }

Set the environment variables and run:

PRIOSTACK_URL=https://priostack.com \ PRIOSTACK_API_KEY=your-api-key \ go run .

Conclusion

The BPMN service task worker pattern in Go is straightforward: poll, execute, complete. The key design decisions are:

  • Pause between empty polls: activation answers at once, and the limit is 300 requests per minute per IP.
  • Run each job in its own goroutine for concurrency, and decode every key as a string.
  • Fail transient errors with job.Retries - 1; report business outcomes as output variables the model routes on.
  • Check the status of every call, and on shutdown wait for jobs in flight, then release the rest with unchanged retries.

Priostack's API keeps things simple - no gRPC, no heavy client library, just HTTP. You can have your first worker running in minutes.

Ready to run your first BPMN service task?

Get a free API key, deploy your BPMN, and have a worker polling in under 5 minutes.

Quickstart guide API docs

Frequently Asked Questions

What is a BPMN service task worker?

A BPMN service task worker is a long-running process that polls a workflow engine for pending jobs of a specific type, executes the business logic, and reports the result back to the engine. The worker pattern decouples your business logic from the workflow orchestration layer.

How do I implement a BPMN worker in Go?

With Priostack, poll for activated jobs via POST /api/v1/jobs/activate, process each job in a goroutine, and complete them via POST /api/v1/jobs/{key}/complete. See the full code example above.

What is the difference between a job worker and a service task?

A service task is the BPMN modelling concept - a task in your process diagram that calls an external service. A job worker is the runtime implementation - the Go process that subscribes to that service task type and executes it.

How does error handling work in BPMN service task workers?

Workers can report job failure via POST /api/v1/jobs/{key}/fail with a decremented retries count. When retries reach zero the engine creates an incident. There is no REST call to throw a BPMN error: complete the job with an output variable and route on it with a gateway.

Can I use Priostack as a Camunda/Zeebe worker replacement?

Partly. Priostack follows the shape of Zeebe's job calls over HTTPS and JSON under /api/v1, with no gRPC, so a Zeebe client library cannot connect. Your workers keep their job logic, but their transport is rewritten to call these endpoints with an API key. See also: migrating from Camunda.

Related: Enterprise integration patterns without a message broker  ·  Two-layer BPMN architecture  ·  Docs: Job Workers