zdravko/pkg/worker/worker.go

108 lines
2.6 KiB
Go
Raw Normal View History

2024-02-16 21:31:00 +00:00
package worker
import (
"encoding/json"
"io"
"log"
"net/http"
2024-02-21 09:41:20 +00:00
"time"
2024-02-16 21:31:00 +00:00
"code.tjo.space/mentos1386/zdravko/internal/activities"
"code.tjo.space/mentos1386/zdravko/internal/config"
"code.tjo.space/mentos1386/zdravko/internal/temporal"
"code.tjo.space/mentos1386/zdravko/internal/workflows"
"code.tjo.space/mentos1386/zdravko/pkg/api"
2024-02-21 09:41:20 +00:00
"code.tjo.space/mentos1386/zdravko/pkg/retry"
"github.com/pkg/errors"
2024-02-16 21:31:00 +00:00
"go.temporal.io/sdk/worker"
)
type ConnectionConfig struct {
Endpoint string `json:"endpoint"`
Group string `json:"group"`
}
func getConnectionConfig(token string, apiUrl string) (*ConnectionConfig, error) {
req, err := api.NewRequest(http.MethodGet, apiUrl+"/api/v1/workers/connect", token, nil)
if err != nil {
return nil, err
}
2024-02-21 09:41:20 +00:00
return retry.Retry(10, 3*time.Second, func() (*ConnectionConfig, error) {
res, err := http.DefaultClient.Do(req)
if err != nil {
return nil, errors.Wrap(err, "failed to connect to API")
}
2024-02-23 11:18:02 +00:00
if res.StatusCode == http.StatusUnauthorized {
2024-02-24 21:07:49 +00:00
panic("WORKER_GROUP_TOKEN is invalid. Either it expired or the worker was removed!")
2024-02-23 11:18:02 +00:00
}
if res.StatusCode != http.StatusOK {
return nil, errors.Errorf("unexpected status code: %d", res.StatusCode)
}
2024-02-21 09:41:20 +00:00
body, err := io.ReadAll(res.Body)
if err != nil {
return nil, errors.Wrap(err, "failed to read response body")
}
config := ConnectionConfig{}
err = json.Unmarshal(body, &config)
if err != nil {
return nil, errors.Wrap(err, "failed to unmarshal connection config")
}
return &config, nil
})
}
2024-02-16 21:31:00 +00:00
type Worker struct {
worker worker.Worker
cfg *config.WorkerConfig
2024-02-16 21:31:00 +00:00
}
func NewWorker(cfg *config.WorkerConfig) (*Worker, error) {
2024-02-16 21:31:00 +00:00
return &Worker{
cfg: cfg,
}, nil
}
func (w *Worker) Name() string {
return "Temporal Worker"
}
func (w *Worker) Start() error {
config, err := getConnectionConfig(w.cfg.Token, w.cfg.ApiUrl)
if err != nil {
return err
}
2024-02-24 21:07:49 +00:00
log.Println("Worker Group:", config.Group)
2024-02-24 21:07:49 +00:00
temporalClient, err := temporal.ConnectWorkerToTemporal(w.cfg.Token, config.Endpoint)
2024-02-16 21:31:00 +00:00
if err != nil {
return err
}
// Create a new Worker
w.worker = worker.New(temporalClient, config.Group, worker.Options{})
2024-02-16 21:31:00 +00:00
workerActivities := activities.NewActivities(w.cfg)
workerWorkflows := workflows.NewWorkflows(workerActivities)
2024-02-16 21:31:00 +00:00
// Register Workflows
w.worker.RegisterWorkflow(workerWorkflows.MonitorWorkflowDefinition)
2024-02-16 21:31:00 +00:00
// Register Activities
w.worker.RegisterActivity(workerActivities.Monitor)
w.worker.RegisterActivity(workerActivities.MonitorAddToHistory)
2024-02-16 21:31:00 +00:00
return w.worker.Run(worker.InterruptCh())
}
func (w *Worker) Stop() error {
w.worker.Stop()
return nil
}