package servicerequest import ( "context" "database/sql" "encoding/json" "fmt" "net/http" "os" "path/filepath" "strconv" "strings" "time" "service/internal/infrastructure/database" extapi "service/internal/worker/interface" "service/pkg/logger" ) const trackerFile = "internal/worker/satusehat/servicerequest/last_servicerequest_id.txt" const trackerFileRad = "internal/worker/satusehat/servicerequest/last_servicerequest_rad_time.txt" const rateLimitSleep = 60 * time.Second type Config struct { DBManager database.Service InternalBaseURL string OrganizationID string } type TokenManager interface { GetAccessToken() string ForceRefreshAndGetToken(ctx context.Context) (string, error) } type WorkerService interface { Run(ctx context.Context) } type worker struct { cfg Config repo Repository apiClient extapi.Client tokenManager TokenManager mode string } func NewDefaultWorker(cfg Config, tokenManager TokenManager) WorkerService { return &worker{cfg: cfg, repo: NewRepository(cfg.DBManager), apiClient: extapi.NewClient(15 * time.Second), tokenManager: tokenManager, mode: "default"} } func NewRadiologyWorker(cfg Config, tokenManager TokenManager) WorkerService { return &worker{cfg: cfg, repo: NewRepository(cfg.DBManager), apiClient: extapi.NewClient(15 * time.Second), tokenManager: tokenManager, mode: "radiology"} } func (w *worker) Run(ctx context.Context) { if w.mode == "radiology" { w.runRadiology(ctx) } else { w.runDefault(ctx) } } func (w *worker) runDefault(ctx context.Context) { logger.Default().Info("[SERVICEREQUEST WORKER] Memulai proses sinkronisasi...") w.ensureTrackerFileExists() lastID := w.readLastID() logger.Default().Info("[SERVICEREQUEST WORKER] Melanjutkan proses dari", logger.Int64("last_id", lastID)) for { select { case <-ctx.Done(): logger.Default().Warn("[SERVICEREQUEST WORKER] Proses dihentikan oleh sistem (Graceful Shutdown)") return default: } sr, err := w.repo.GetNextServiceRequest(ctx, lastID) if err != nil { if err == sql.ErrNoRows { time.Sleep(5 * time.Second) // Tunggu data baru continue } logger.Default().Error("[SERVICEREQUEST WORKER] Gagal query ke database", logger.ErrorField(err)) time.Sleep(5 * time.Second) continue } sourceID := strconv.FormatInt(sr.IdxOrder, 10) logger.Default().Info("[SERVICEREQUEST WORKER] 🔄 Memproses Dokumen", logger.String("source_id", sourceID), logger.String("no_mr", sr.NoMR.String), ) internalPayload := MapToInternalAPI(sr, w.cfg.OrganizationID) patientID, _ := internalPayload["patient_id"].(string) encounterID, _ := internalPayload["encounter_id"].(string) if patientID == "" || encounterID == "" || encounterID == "0" { logger.Default().Warn("[SERVICEREQUEST WORKER] ⚠️ Data dilewati karena patient_id atau encounter_id kosong", logger.String("source_id", sourceID)) lastID = sr.IdxOrder os.WriteFile(trackerFile, []byte(strconv.FormatInt(lastID, 10)), 0644) continue } fhirID, respBody, status, errMsg, retryDelay := w.sendToInternalAPI(ctx, internalPayload, sourceID) syncLog := MapToSyncLog(sr.IdxOrder, fhirID, internalPayload, respBody, status, errMsg) if errLog := w.repo.SaveSyncLog(ctx, syncLog); errLog != nil { logger.Default().Error("[SERVICEREQUEST WORKER] Gagal menyimpan histori log ke DB", logger.ErrorField(errLog)) } if status == "SUCCESS" && fhirID != "" { if errUpdate := w.repo.SaveSatuSehatID(ctx, sr.IdxOrder, fhirID); errUpdate != nil { logger.Default().Error("[SERVICEREQUEST WORKER] Gagal update FHIR ID ke transaksi utama", logger.ErrorField(errUpdate)) } } if retryDelay == 0 { lastID = sr.IdxOrder os.WriteFile(trackerFile, []byte(strconv.FormatInt(lastID, 10)), 0644) } else { logger.Default().Warn("[SERVICEREQUEST WORKER] ⚠️ Retrying data", logger.String("source_id", sourceID), logger.String("reason", errMsg), logger.Duration("delay", retryDelay)) } if retryDelay > 0 { time.Sleep(retryDelay) } else { time.Sleep(2 * time.Second) // Mencegah Rate Limit Exceeded } } } func (w *worker) runRadiology(ctx context.Context) { logger.Default().Info("[SERVICEREQUEST-RAD WORKER] Memulai proses sinkronisasi radiologi...") w.ensureTrackerFileRadExists() lastTime := w.readLastTime() // Custom data: Override membaca mulai tanggal 2025-12-18 (Format: YYYY-MM-DD) customDate, _ := time.Parse("2006-01-02", "2026-04-24") if lastTime.Before(customDate) { lastTime = customDate } logger.Default().Info("[SERVICEREQUEST-RAD WORKER] Melanjutkan proses dari", logger.Time("last_time", lastTime)) for { select { case <-ctx.Done(): logger.Default().Warn("[SERVICEREQUEST-RAD WORKER] Proses dihentikan oleh sistem (Graceful Shutdown)") return default: } sr, err := w.repo.GetNextRadiologyServiceRequest(ctx, lastTime) if err != nil { if err == sql.ErrNoRows { time.Sleep(5 * time.Second) // Tunggu data baru continue } logger.Default().Error("[SERVICEREQUEST-RAD WORKER] Gagal query ke database", logger.ErrorField(err)) time.Sleep(5 * time.Second) continue } logger.Default().Info("[SERVICEREQUEST-RAD WORKER] 🔄 Memproses Dokumen", logger.String("source_id", sr.ID), logger.String("patient_name", sr.PatientName.String), ) internalPayload := MapRadiologyToInternalAPI(sr) patientID, _ := internalPayload["patient_id"].(string) encounterID, _ := internalPayload["encounter_id"].(string) requesterID, _ := internalPayload["requester_id"].(string) performerID, _ := internalPayload["performer_id"].(string) if patientID == "" || encounterID == "" || encounterID == "0" || requesterID == "" || requesterID == "0" || performerID == "" || performerID == "0" { logger.Default().Warn("[SERVICEREQUEST-RAD WORKER] ⚠️ Data dilewati karena patient_id atau encounter_id atau requester_id kosong", logger.String("source_id", sr.ID)) lastTime = sr.CreatedAt os.WriteFile(trackerFileRad, []byte(lastTime.Format(time.RFC3339Nano)), 0644) continue } fhirID, respBody, status, errMsg, retryDelay := w.sendToInternalAPI(ctx, internalPayload, sr.ID) // Menyimpan payload request dan response ke dalam log logger.Default().Info("[SERVICEREQUEST-RAD WORKER] Detail Transaksi", logger.String("source_id", sr.ID), logger.Any("request_payload", internalPayload), logger.String("response_payload", string(respBody)), logger.String("status", status), ) if status == "SUCCESS" && fhirID != "" { if errUpdate := w.repo.UpdateRadiologyServiceRequestStatus(ctx, sr.ID, fhirID, encounterID, true); errUpdate != nil { logger.Default().Error("[SERVICEREQUEST-RAD WORKER] Gagal update status ke transaksi utama", logger.ErrorField(errUpdate)) } } if retryDelay == 0 { lastTime = sr.CreatedAt os.WriteFile(trackerFileRad, []byte(lastTime.Format(time.RFC3339Nano)), 0644) } else { logger.Default().Warn("[SERVICEREQUEST-RAD WORKER] ⚠️ Retrying data", logger.String("source_id", sr.ID), logger.String("reason", errMsg), logger.Duration("delay", retryDelay)) } if retryDelay > 0 { time.Sleep(retryDelay) } else { time.Sleep(2 * time.Second) // Mencegah Rate Limit Exceeded } } } // isRateLimitResponse mendeteksi respons rate-limit dari Satu Sehat // (HTTP 429 atau pesan body yang mengindikasikan limit pengiriman terlampaui). func isRateLimitResponse(statusCode int, respBody []byte) bool { if statusCode == http.StatusTooManyRequests { return true } body := strings.ToLower(string(respBody)) return strings.Contains(body, "too many requests") || strings.Contains(body, "rate limit") || strings.Contains(body, "limit exceeded") || strings.Contains(body, "limit pengiriman") || strings.Contains(body, "quota exceeded") } // sendToInternalAPI mengirimkan payload ke internal FHIR server, menangani auth retry (401) dan response error. func (w *worker) sendToInternalAPI(ctx context.Context, payload map[string]interface{}, sourceID string) (fhirID string, respBody []byte, status string, errMsg string, retryDelay time.Duration) { reqURL := fmt.Sprintf("%s/satusehat/servicerequest", strings.TrimRight(w.cfg.InternalBaseURL, "/")) workerName := "SERVICEREQUEST WORKER" if w.mode == "radiology" { workerName = "SERVICEREQUEST-RAD WORKER" } // --- Percobaan Pertama --- respBody, statusCode, err := w.apiClient.PostJSON(ctx, reqURL, w.tokenManager.GetAccessToken(), payload) if err != nil { logger.Default().Error(fmt.Sprintf("[%s] ❌ Gagal Koneksi ke API Part 3", workerName), logger.String("source_id", sourceID), logger.Any("payload_sent", payload), logger.ErrorField(err), ) return "", nil, "FAILED", err.Error(), 5 * time.Second } // --- Penanganan Token Unauthorized (401) --- if statusCode == http.StatusUnauthorized { logger.Default().Warn(fmt.Sprintf("[%s] Menerima status 401 Unauthorized, mencoba refresh token...", workerName), logger.String("source_id", sourceID)) newToken, refreshErr := w.tokenManager.ForceRefreshAndGetToken(ctx) if refreshErr != nil { logger.Default().Error(fmt.Sprintf("[%s] Gagal refresh token, akan mencoba lagi pada iterasi berikutnya.", workerName), logger.String("source_id", sourceID), logger.ErrorField(refreshErr)) return "", respBody, "FAILED", refreshErr.Error(), 5 * time.Second } logger.Default().Info(fmt.Sprintf("[%s] Token berhasil di-refresh, mencoba ulang request...", workerName), logger.String("source_id", sourceID)) respBody, statusCode, err = w.apiClient.PostJSON(ctx, reqURL, newToken, payload) if err != nil { logger.Default().Error(fmt.Sprintf("[%s] ❌ Gagal Koneksi ke API Part 3 (setelah retry)", workerName), logger.String("source_id", sourceID), logger.ErrorField(err), ) return "", nil, "FAILED", err.Error(), 5 * time.Second } } // --- Evaluasi Final Respons --- if statusCode != http.StatusOK && statusCode != http.StatusCreated { errMsg = fmt.Sprintf("HTTP %d", statusCode) logger.Default().Error(fmt.Sprintf("[%s] ❌ Respons Error dari API Part 3", workerName), logger.String("source_id", sourceID), logger.Int("status_code", statusCode), logger.String("response", truncate(string(respBody), 200)), ) // Lakukan Retry jika server Internal error (5xx) atau Rate Limit terdeteksi if statusCode >= 500 || isRateLimitResponse(statusCode, respBody) { delay := 5 * time.Second // Default delay untuk 5xx if isRateLimitResponse(statusCode, respBody) { delay = rateLimitSleep // Gunakan konstanta dari reference var errResp struct { RetryAfter int `json:"retry_after"` } if errParse := json.Unmarshal(respBody, &errResp); errParse == nil && errResp.RetryAfter > 0 { delay = time.Duration(errResp.RetryAfter) * time.Second } } return "", respBody, "FAILED", errMsg, delay } return "", respBody, "FAILED", errMsg, 0 } // --- Sukses --- var result map[string]interface{} if errUnmarshal := json.Unmarshal(respBody, &result); errUnmarshal == nil { if data, ok := result["data"].(map[string]interface{}); ok { if id, ok := data["id"].(string); ok { fhirID = id } } else if id, ok := result["id"].(string); ok { fhirID = id } } logger.Default().Info(fmt.Sprintf("[%s] ✅ Sukses mengirim data ServiceRequest", workerName), logger.String("source_id", sourceID), logger.String("fhir_id", fhirID), ) return fhirID, respBody, "SUCCESS", "", 0 } func (w *worker) readLastID() int64 { data, _ := os.ReadFile(trackerFile) id, _ := strconv.ParseInt(strings.TrimSpace(string(data)), 10, 64) return id } func (w *worker) ensureTrackerFileExists() { if _, err := os.Stat(trackerFile); os.IsNotExist(err) { dir := filepath.Dir(trackerFile) if err := os.MkdirAll(dir, 0755); err != nil { logger.Default().Error("[SERVICEREQUEST WORKER] Gagal membuat direktori tracker", logger.ErrorField(err)) } if err := os.WriteFile(trackerFile, []byte("0"), 0644); err != nil { logger.Default().Error("[SERVICEREQUEST WORKER] Gagal membuat file tracker default", logger.ErrorField(err)) } } } func (w *worker) readLastTime() time.Time { data, err := os.ReadFile(trackerFileRad) if err != nil { return time.Time{} } t, err := time.Parse(time.RFC3339Nano, strings.TrimSpace(string(data))) if err != nil { return time.Time{} } return t } func (w *worker) ensureTrackerFileRadExists() { if _, err := os.Stat(trackerFileRad); os.IsNotExist(err) { dir := filepath.Dir(trackerFileRad) _ = os.MkdirAll(dir, 0755) _ = os.WriteFile(trackerFileRad, []byte(time.Time{}.Format(time.RFC3339Nano)), 0644) } } func truncate(s string, max int) string { if len(s) > max { return s[:max] + "..." } return s }