update ignore
This commit is contained in:
No files matched your search
@@ -0,0 +1,331 @@
|
||||
package worker
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"os"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"service/internal/infrastructure/config"
|
||||
"service/internal/infrastructure/database"
|
||||
"service/internal/interfaces/satusehat"
|
||||
"service/internal/master/kfa"
|
||||
"service/pkg/logger"
|
||||
)
|
||||
|
||||
// Manager mengelola siklus hidup semua background workers
|
||||
type Manager struct {
|
||||
cfg *config.Config
|
||||
db database.Service
|
||||
satusehatClient satusehat.SatuSehatClient
|
||||
|
||||
// Auth state for all workers
|
||||
accessToken string
|
||||
refreshToken string
|
||||
tokenMutex sync.RWMutex
|
||||
}
|
||||
|
||||
// NewManager membuat instance Manager baru
|
||||
func NewManager(cfg *config.Config, db database.Service, satusehatClient satusehat.SatuSehatClient) *Manager {
|
||||
return &Manager{
|
||||
cfg: cfg,
|
||||
db: db,
|
||||
satusehatClient: satusehatClient,
|
||||
}
|
||||
}
|
||||
|
||||
// Start menjalankan semua worker terdaftar dalam goroutine terpisah
|
||||
func (m *Manager) Start(ctx context.Context) {
|
||||
logger.Default().Info("Starting background workers...")
|
||||
|
||||
// Lakukan login awal untuk mendapatkan token bagi semua worker
|
||||
maxRetries := 6
|
||||
for i := 1; i <= maxRetries; i++ {
|
||||
err := m.login(ctx)
|
||||
if err == nil {
|
||||
break
|
||||
}
|
||||
logger.Default().Warn(fmt.Sprintf("Initial login attempt %d/%d failed, retrying in 5 seconds...", i, maxRetries), logger.ErrorField(err))
|
||||
if i == maxRetries {
|
||||
logger.Default().Fatal("Initial login for workers failed after maximum retries, stopping.", logger.ErrorField(err))
|
||||
return
|
||||
}
|
||||
time.Sleep(5 * time.Second)
|
||||
}
|
||||
|
||||
var wg sync.WaitGroup
|
||||
|
||||
// Tentukan Organization ID (Bisa dari Environment Variable atau Config Struct)
|
||||
orgID := os.Getenv("SATUSEHAT_ORG_ID")
|
||||
if orgID == "" {
|
||||
logger.Default().Warn("SATUSEHAT_ORG_ID is not set in environment variables")
|
||||
}
|
||||
|
||||
internalBaseURL := m.cfg.SatuSehat.InternalFHIRServerURL
|
||||
|
||||
// Jobs mendefinisikan urutan dan jeda waktu (delay) boot-up setiap worker.
|
||||
// Konfigurasi waktu ini mencegah perebutan resource di awal (thundering herd)
|
||||
// dan memastikan resource inti (seperti Encounter) mendapat "head start"
|
||||
// sebelum worker anak mencoba merujuk ke data tersebut.
|
||||
jobs := []struct {
|
||||
name string
|
||||
delay time.Duration
|
||||
run func(context.Context)
|
||||
}{
|
||||
// {
|
||||
// name: "Encounter",
|
||||
// delay: 5 * time.Second, // Memberikan jeda 5 detik agar Encounter mendapat porsi jalan lebih dulu
|
||||
// run: func(c context.Context) {
|
||||
// encounter.NewWorker(encounter.Config{
|
||||
// DBManager: m.db, InternalBaseURL: internalBaseURL, OrganizationID: orgID,
|
||||
// }).Run(c)
|
||||
// },
|
||||
// },
|
||||
// {
|
||||
// name: "Condition",
|
||||
// delay: 2 * time.Second,
|
||||
// run: func(c context.Context) {
|
||||
// // TODO: Refactor this worker
|
||||
// // condition.NewWorker(...).Run(c)
|
||||
// },
|
||||
// },
|
||||
// {
|
||||
// name: "Medication (Master)",
|
||||
// delay: 2, // Master data berjalan seketika
|
||||
// run: func(c context.Context) {
|
||||
// medication.NewWorker(medication.Config{
|
||||
// DBManager: m.db,
|
||||
// InternalBaseURL: internalBaseURL,
|
||||
// OrganizationID: orgID,
|
||||
// }, m).Run(c)
|
||||
// },
|
||||
// },
|
||||
// {
|
||||
// name: "Observation",
|
||||
// delay: 2 * time.Second,
|
||||
// run: func(c context.Context) {
|
||||
// // TODO: Refactor this worker
|
||||
// // observation.NewWorker(...).Run(c)
|
||||
// },
|
||||
// },
|
||||
// {
|
||||
// name: "Procedure",
|
||||
// delay: 2 * time.Second,
|
||||
// run: func(c context.Context) {
|
||||
// // TODO: Refactor this worker
|
||||
// // procedure.NewWorker(...).Run(c)
|
||||
// },
|
||||
// },
|
||||
// {
|
||||
// name: "AllergyIntolerance",
|
||||
// delay: 2 * time.Second,
|
||||
// run: func(c context.Context) {
|
||||
// // TODO: Refactor this worker
|
||||
// // allergyintolerance.NewWorker(...).Run(c)
|
||||
// },
|
||||
// },
|
||||
// {
|
||||
// name: "ClinicalImpression",
|
||||
// delay: 2 * time.Second,
|
||||
// run: func(c context.Context) {
|
||||
// // TODO: Refactor this worker
|
||||
// // clinicalimpression.NewWorker(...).Run(c)
|
||||
// },
|
||||
// },
|
||||
// {
|
||||
// name: "Immunization",
|
||||
// delay: 2 * time.Second,
|
||||
// run: func(c context.Context) {
|
||||
// // TODO: Refactor this worker
|
||||
// // immunization.NewWorker(...).Run(c)
|
||||
// },
|
||||
// },
|
||||
// {
|
||||
// name: "ServiceRequest (Default)",
|
||||
// delay: 0,
|
||||
// run: func(c context.Context) {
|
||||
// servicerequest.NewDefaultWorker(servicerequest.Config{
|
||||
// DBManager: m.db,
|
||||
// InternalBaseURL: internalBaseURL,
|
||||
// OrganizationID: orgID,
|
||||
// }, m).Run(c)
|
||||
// },
|
||||
// },
|
||||
// {
|
||||
// name: "ServiceRequest (Radiology)",
|
||||
// delay: 2 * time.Second,
|
||||
// run: func(c context.Context) {
|
||||
// servicerequest.NewRadiologyWorker(servicerequest.Config{
|
||||
// DBManager: m.db,
|
||||
// InternalBaseURL: internalBaseURL,
|
||||
// OrganizationID: orgID,
|
||||
// }, m).Run(c)
|
||||
// },
|
||||
// },
|
||||
// {
|
||||
// name: "Specimen",
|
||||
// delay: 2 * time.Second,
|
||||
// run: func(c context.Context) {
|
||||
// // TODO: Refactor this worker
|
||||
// // specimen.NewWorker(...).Run(c)
|
||||
// },
|
||||
// },
|
||||
// {
|
||||
// name: "DiagnosticReport",
|
||||
// delay: 2 * time.Second,
|
||||
// run: func(c context.Context) {
|
||||
// // TODO: Refactor this worker
|
||||
// // diagnosticreport.NewWorker(...).Run(c)
|
||||
// },
|
||||
// },
|
||||
// {
|
||||
// name: "ImagingStudy",
|
||||
// delay: 2 * time.Second,
|
||||
// run: func(c context.Context) {
|
||||
// // The Manager 'm' now acts as the TokenManager
|
||||
// imagingstudy.NewWorker(imagingstudy.Config{
|
||||
// DBManager: m.db, InternalBaseURL: internalBaseURL, OrganizationID: orgID,
|
||||
// }, m).Run(c)
|
||||
// },
|
||||
// },
|
||||
// {
|
||||
// name: "Composition",
|
||||
// delay: 2 * time.Second,
|
||||
// run: func(c context.Context) {
|
||||
// // TODO: Refactor this worker
|
||||
// // composition.NewWorker(...).Run(c)
|
||||
// },
|
||||
// },
|
||||
// {
|
||||
// name: "CarePlan",
|
||||
// delay: 2 * time.Second,
|
||||
// run: func(c context.Context) {
|
||||
// // TODO: Refactor this worker
|
||||
// // careplan.NewWorker(...).Run(c)
|
||||
// },
|
||||
// },
|
||||
// {
|
||||
// name: "EpisodeOfCare",
|
||||
// delay: 2 * time.Second,
|
||||
// run: func(c context.Context) {
|
||||
// // TODO: Refactor this worker
|
||||
// // episodeofcare.NewWorker(...).Run(c)
|
||||
// },
|
||||
// },
|
||||
// {
|
||||
// name: "MedicationRequest",
|
||||
// delay: 2 * time.Second,
|
||||
// run: func(c context.Context) {
|
||||
// medicationrequest.NewWorker(medicationrequest.Config{
|
||||
// DBManager: m.db,
|
||||
// InternalBaseURL: internalBaseURL,
|
||||
// OrganizationID: orgID,
|
||||
// }, m).Run(c)
|
||||
// },
|
||||
// },
|
||||
// {
|
||||
// name: "MedicationDispense",
|
||||
// delay: 2 * time.Second,
|
||||
// run: func(c context.Context) {
|
||||
// medicationdispense.NewWorker(medicationdispense.Config{
|
||||
// DBManager: m.db,
|
||||
// InternalBaseURL: internalBaseURL,
|
||||
// OrganizationID: orgID,
|
||||
// }, m).Run(c)
|
||||
// },
|
||||
// },
|
||||
{
|
||||
name: "KFA Master Puller",
|
||||
delay: 5 * time.Second,
|
||||
run: func(c context.Context) {
|
||||
repo := kfa.NewCommandRepository(m.db, "default")
|
||||
kfa.NewWorker(kfa.Config{
|
||||
InternalBaseURL: internalBaseURL,
|
||||
PageSize: 100,
|
||||
}, repo, m).Run(c)
|
||||
},
|
||||
},
|
||||
// {
|
||||
// name: "MedicationStatement",
|
||||
// delay: 2 * time.Second,
|
||||
// run: func(c context.Context) {
|
||||
// // TODO: Refactor this worker
|
||||
// // medicationstatement.NewWorker(...).Run(c)
|
||||
// },
|
||||
// },
|
||||
// {
|
||||
// name: "QuestionnaireResponse",
|
||||
// delay: 2 * time.Second,
|
||||
// run: func(c context.Context) {
|
||||
// // TODO: Refactor this worker
|
||||
// // questionnaireresponse.NewWorker(...).Run(c)
|
||||
// },
|
||||
// },
|
||||
}
|
||||
|
||||
// Eksekusi semua background worker secara tertib dengan delay bertingkat
|
||||
go func() {
|
||||
baseDelay := 10 * time.Second // Jeda 10 detik antar tiap worker untuk cegah spike di awal
|
||||
|
||||
for i, job := range jobs {
|
||||
// Hitung delay bertingkat: (index) * baseDelay
|
||||
staggeredDelay := time.Duration(i) * baseDelay
|
||||
|
||||
// Jika job memiliki delay statis yang lebih besar, utamakan nilai tersebut
|
||||
if job.delay > staggeredDelay {
|
||||
staggeredDelay = job.delay
|
||||
}
|
||||
|
||||
wg.Add(1)
|
||||
j := job
|
||||
d := staggeredDelay
|
||||
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
|
||||
if d > 0 {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case <-time.After(d):
|
||||
}
|
||||
}
|
||||
|
||||
logger.Default().Info(fmt.Sprintf("[%s WORKER] Started", j.name))
|
||||
j.run(ctx)
|
||||
logger.Default().Info(fmt.Sprintf("[%s WORKER] Stopped", j.name))
|
||||
}()
|
||||
}
|
||||
|
||||
// Waiter untuk Graceful Shutdown keseluruhan
|
||||
wg.Wait()
|
||||
logger.Default().Info("All background workers stopped gracefully")
|
||||
}()
|
||||
}
|
||||
|
||||
// StartKFAWorker menjalankan worker khusus untuk menarik master data KFA secara independen
|
||||
func (m *Manager) StartKFAWorker(ctx context.Context) {
|
||||
logger.Default().Info("Starting independent KFA Master Puller worker...")
|
||||
|
||||
// Lakukan autentikasi awal jika token belum tersedia (Aman untuk dipanggil secara konkuren)
|
||||
if m.GetAccessToken() == "" {
|
||||
if _, err := m.ForceRefreshAndGetToken(ctx); err != nil {
|
||||
logger.Default().Fatal("Initial login for KFA worker failed, stopping.", logger.ErrorField(err))
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
internalBaseURL := m.cfg.SatuSehat.InternalFHIRServerURL
|
||||
|
||||
// Eksekusi KFA worker di goroutine yang terpisah dari pool antrean Manager.Start()
|
||||
go func() {
|
||||
logger.Default().Info("[KFA Master Puller WORKER] Started")
|
||||
repo := kfa.NewCommandRepository(m.db, "default")
|
||||
kfa.NewWorker(kfa.Config{
|
||||
InternalBaseURL: internalBaseURL,
|
||||
PageSize: 100,
|
||||
}, repo, m).Run(ctx)
|
||||
logger.Default().Info("[KFA Master Puller WORKER] Stopped")
|
||||
}()
|
||||
}
|
||||
Reference in New Issue
Block a user