Add Temporal worker and test client programs
Creates: - cmd/worker/main.go: Worker that registers workflows and activities - cmd/test-workflow/main.go: Test client to trigger workflows Adds go.temporal.io/sdk dependency to go.mod.
This commit is contained in:
@@ -0,0 +1,120 @@
|
||||
package workflow
|
||||
|
||||
import (
|
||||
"log"
|
||||
"os"
|
||||
"os/signal"
|
||||
"syscall"
|
||||
"time"
|
||||
|
||||
"go.temporal.io/sdk/client"
|
||||
"go.temporal.io/sdk/worker"
|
||||
"go.temporal.io/sdk/temporal"
|
||||
"go.temporal.io/sdk/workflow"
|
||||
)
|
||||
|
||||
// WorkerConfig holds worker configuration
|
||||
type WorkerConfig struct {
|
||||
HostPort string
|
||||
Namespace string
|
||||
TaskQueue string
|
||||
}
|
||||
|
||||
// NewWorker creates and starts a Temporal worker
|
||||
func NewWorker(cfg WorkerConfig) error {
|
||||
// Connect to Temporal server
|
||||
c, err := client.Dial(client.Options{
|
||||
HostPort: cfg.HostPort,
|
||||
Namespace: cfg.Namespace,
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer c.Close()
|
||||
|
||||
log.Printf("Connected to Temporal at %s (namespace: %s)", cfg.HostPort, cfg.Namespace)
|
||||
|
||||
// Create worker
|
||||
w := worker.New(c, cfg.TaskQueue, worker.Options{})
|
||||
|
||||
// Register workflows
|
||||
w.RegisterWorkflow(HelloWorldWorkflow)
|
||||
w.RegisterWorkflow(GreeterWorkflow)
|
||||
w.RegisterWorkflow(ProcessOrderWorkflow)
|
||||
|
||||
// Register activities
|
||||
w.RegisterActivity(GreetActivity)
|
||||
w.RegisterActivity(ValidateOrderActivity)
|
||||
w.RegisterActivity(ProcessPaymentActivity)
|
||||
w.RegisterActivity(NotifyCustomerActivity)
|
||||
|
||||
// Start worker (blocks until signal received)
|
||||
log.Printf("Starting worker on task queue: %s", cfg.TaskQueue)
|
||||
|
||||
sigChan := make(chan os.Signal, 1)
|
||||
signal.Notify(sigChan, syscall.SIGINT, syscall.SIGTERM)
|
||||
|
||||
go func() {
|
||||
if err := w.Run(worker.InterruptCh()); err != nil {
|
||||
log.Fatal(err)
|
||||
}
|
||||
}()
|
||||
|
||||
// Wait for shutdown signal
|
||||
<-sigChan
|
||||
log.Println("Shutting down worker...")
|
||||
w.Stop()
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// ProcessOrderWorkflow demonstrates multi-step workflow with activities
|
||||
func ProcessOrderWorkflow(ctx workflow.Context, orderID string) (string, error) {
|
||||
opts := workflow.ActivityOptions{
|
||||
StartToCloseTimeout: 5 * time.Minute,
|
||||
RetryPolicy: &temporal.RetryPolicy{
|
||||
InitialInterval: time.Second,
|
||||
BackoffCoefficient: 2.0,
|
||||
MaximumInterval: time.Minute,
|
||||
MaximumAttempts: 3,
|
||||
},
|
||||
}
|
||||
ctx = workflow.WithActivityOptions(ctx, opts)
|
||||
|
||||
// Step 1: Validate order
|
||||
var validated bool
|
||||
if err := workflow.ExecuteActivity(ctx, ValidateOrderActivity, orderID).Get(ctx, &validated); err != nil {
|
||||
return "", err
|
||||
}
|
||||
if !validated {
|
||||
return "", temporal.NewApplicationError("invalid order", "InvalidOrder")
|
||||
}
|
||||
|
||||
// Step 2: Process payment
|
||||
var paymentID string
|
||||
if err := workflow.ExecuteActivity(ctx, ProcessPaymentActivity, orderID).Get(ctx, &paymentID); err != nil {
|
||||
return "", err
|
||||
}
|
||||
|
||||
// Step 3: Notify customer
|
||||
var notifyResult string
|
||||
if err := workflow.ExecuteActivity(ctx, NotifyCustomerActivity, orderID).Get(ctx, ¬ifyResult); err != nil {
|
||||
return "", err
|
||||
}
|
||||
|
||||
return paymentID, nil
|
||||
}
|
||||
|
||||
// GreeterWorkflow is a multi-step workflow
|
||||
func GreeterWorkflow(ctx workflow.Context, name string) (string, error) {
|
||||
opts := workflow.ActivityOptions{
|
||||
StartToCloseTimeout: 5 * time.Minute,
|
||||
}
|
||||
ctx = workflow.WithActivityOptions(ctx, opts)
|
||||
|
||||
var result string
|
||||
if err := workflow.ExecuteActivity(ctx, GreetActivity, name).Get(ctx, &result); err != nil {
|
||||
return "", err
|
||||
}
|
||||
return result, nil
|
||||
}
|
||||
Reference in New Issue
Block a user