package imagingstudy import ( "context" "database/sql" "encoding/json" "fmt" "net/http" "os" "path/filepath" "strings" "time" "service/internal/infrastructure/database" extapi "service/internal/worker/interface" "service/pkg/logger" ) const ( trackerFile = "internal/worker/satusehat/imagingstudy/last_imagingstudy_id.txt" rateLimitSleep = 60 * time.Second ) type Config struct { DBManager database.Service InternalBaseURL string OrganizationID string } // TokenManager defines the interface for a centralized token provider. 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 } func NewWorker(cfg Config, tokenManager TokenManager) WorkerService { return &worker{ cfg: cfg, repo: NewRepository(cfg.DBManager), apiClient: extapi.NewClient(60 * time.Second), tokenManager: tokenManager, } } func (w *worker) Run(ctx context.Context) { logger.Default().Info("[IMAGING STUDY WORKER] Memulai proses sinkronisasi API...") lastCreatedAt := w.readLastCreatedAt() logger.Default().Info("[IMAGING STUDY WORKER] Melanjutkan proses dari", logger.Time("last_created_at", lastCreatedAt)) for { select { case <-ctx.Done(): logger.Default().Warn("[IMAGING STUDY WORKER] Proses dihentikan oleh sistem (Graceful Shutdown)") return default: } // 1. Fetch Data imagingData, err := w.repo.GetNextImagingStudy(ctx, lastCreatedAt) if err != nil { if err == sql.ErrNoRows { time.Sleep(5 * time.Second) // Tunggu data baru continue } logger.Default().Error("[IMAGING STUDY WORKER] Gagal query ke database", logger.ErrorField(err)) time.Sleep(5 * time.Second) continue } // Fallback pencarian Encounter ID lewat API jika data di tabel kosong if (!imagingData.EncounterIHS.Valid || imagingData.EncounterIHS.String == "") && imagingData.ServiceRequestID.Valid && imagingData.ServiceRequestID.String != "" { logger.Default().Info("[IMAGING STUDY WORKER] Encounter ID kosong, mencoba fetch dari API ServiceRequest", logger.String("servicerequest_id", imagingData.ServiceRequestID.String)) reqURL := fmt.Sprintf("%s/satusehat/servicerequest/%s", strings.TrimRight(w.cfg.InternalBaseURL, "/"), imagingData.ServiceRequestID.String) srRespBody, srStatusCode, srErr := w.apiClient.GetJSON(ctx, reqURL, w.tokenManager.GetAccessToken()) if srStatusCode == http.StatusUnauthorized { logger.Default().Warn("[IMAGING STUDY WORKER] 401 Unauthorized saat GET ServiceRequest, refresh token...") if newToken, errRef := w.tokenManager.ForceRefreshAndGetToken(ctx); errRef == nil { srRespBody, srStatusCode, srErr = w.apiClient.GetJSON(ctx, reqURL, newToken) } } if srErr == nil && (srStatusCode == http.StatusOK || srStatusCode == http.StatusCreated) { var parsedResp map[string]interface{} if errParse := json.Unmarshal(srRespBody, &parsedResp); errParse == nil { var encounterRef string if data, ok := parsedResp["data"].(map[string]interface{}); ok { if enc, ok := data["encounter"].(map[string]interface{}); ok { encounterRef, _ = enc["reference"].(string) } } else if enc, ok := parsedResp["encounter"].(map[string]interface{}); ok { encounterRef, _ = enc["reference"].(string) } if encounterRef != "" { if parts := strings.Split(encounterRef, "/"); len(parts) > 1 { imagingData.EncounterIHS.String = parts[len(parts)-1] imagingData.EncounterIHS.Valid = true logger.Default().Info("[IMAGING STUDY WORKER] Berhasil mendapatkan Encounter ID dari API", logger.String("encounter_id", imagingData.EncounterIHS.String)) } } } } } // Gunakan SourceID (UUID) untuk logging dan tracking logger.Default().Info("[IMAGING STUDY WORKER] 🔄 Memproses Dokumen", logger.String("source_id", imagingData.SourceID), logger.String("no_mr", imagingData.NoMR.String), ) // Skip jika accession_number kosong / null if !imagingData.AccessionNumber.Valid || imagingData.AccessionNumber.String == "" { logger.Default().Warn("[IMAGING STUDY WORKER] ⏭️ Skip Dokumen: Accession Number kosong (Data RIS tidak ditemukan)", logger.String("source_id", imagingData.SourceID), logger.String("no_mr", imagingData.NoMR.String), ) w.updateTracker(&lastCreatedAt, imagingData.CreatedAt) continue } // 2. Mapping Payload internalPayload := MapToInternalAPI(imagingData, w.cfg.OrganizationID) reqURL := fmt.Sprintf("%s/satusehat/imagingstudy", strings.TrimRight(w.cfg.InternalBaseURL, "/")) logger.Default().Debug("[IMAGING STUDY WORKER] 📡 Mengirim Request ke API Part 3", logger.String("url", reqURL), logger.Any("payload", internalPayload), ) // 3. Eksekusi Kirim HTTP via Client Helper Anda fhirID, respBody, status, errMsg, retryDelay := w.sendToInternalAPI(ctx, internalPayload, imagingData.SourceID) // 4. Catat Log ke Database syncLog := MapToSyncLog(imagingData.SourceID, fhirID, internalPayload, respBody, status, errMsg) if errLog := w.repo.SaveSyncLog(ctx, syncLog); errLog != nil { logger.Default().Error("[IMAGING STUDY WORKER] Gagal menyimpan histori log ke DB", logger.ErrorField(errLog)) } // 5. Update FHIR ID di tabel utama bila sukses if status == "SUCCESS" && fhirID != "" { if errUpdate := w.repo.SaveSatuSehatID(ctx, imagingData.SourceID, fhirID); errUpdate != nil { logger.Default().Error("[IMAGING STUDY WORKER] Gagal update FHIR ID ke transaksi utama", logger.ErrorField(errUpdate)) } // Tambahkan mapping untuk digunakan oleh DiagnosticReport worker w.appendImagingStudyMapping( imagingData.CreatedAt, imagingData.SourceID, fhirID, imagingData.ServiceRequestID.String, imagingData.EncounterIHS.String, imagingData.PatientIHS.String, imagingData.PractitionerIHS.String, ) } // Update Tracker TXT agar tidak dikirim ulang terus-menerus KECUALI jika butuh retry if retryDelay == 0 { w.updateTracker(&lastCreatedAt, imagingData.CreatedAt) } else { logger.Default().Warn("[IMAGING STUDY WORKER] ⚠️ Retrying data", logger.String("source_id", imagingData.SourceID), logger.String("reason", errMsg), logger.Duration("delay", retryDelay)) } // Jeda if retryDelay > 0 { time.Sleep(retryDelay) } else { time.Sleep(2 * time.Second) // Jeda antar pengiriman diperpanjang untuk mencegah Rate Limit } } } // 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 sends the payload to the internal FHIR server, handles authentication and retries. 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/imagingstudy", strings.TrimRight(w.cfg.InternalBaseURL, "/")) // --- First Attempt --- respBody, statusCode, err := w.apiClient.PostJSON(ctx, reqURL, w.tokenManager.GetAccessToken(), payload) if err != nil { // Network errors are always retried logger.Default().Error("[IMAGING STUDY WORKER] ❌ Gagal Koneksi ke API Part 3", logger.String("source_id", sourceID), logger.Any("payload_sent", payload), logger.ErrorField(err), ) return "", nil, "FAILED", err.Error(), 5 * time.Second } // --- Handle Unauthorized (401) --- if statusCode == http.StatusUnauthorized { logger.Default().Warn("[IMAGING STUDY WORKER] Menerima status 401 Unauthorized, mencoba refresh token...", logger.String("source_id", sourceID)) newToken, refreshErr := w.tokenManager.ForceRefreshAndGetToken(ctx) if refreshErr != nil { logger.Default().Error("[IMAGING STUDY WORKER] Gagal refresh token, akan mencoba lagi pada iterasi berikutnya.", logger.String("source_id", sourceID), logger.ErrorField(refreshErr)) return "", respBody, "FAILED", refreshErr.Error(), 5 * time.Second } logger.Default().Info("[IMAGING STUDY WORKER] Token berhasil di-refresh, mencoba ulang request...", logger.String("source_id", sourceID)) // --- Second Attempt --- respBody, statusCode, err = w.apiClient.PostJSON(ctx, reqURL, newToken, payload) if err != nil { logger.Default().Error("[IMAGING STUDY WORKER] ❌ Gagal Koneksi ke API Part 3 (setelah retry)", logger.String("source_id", sourceID), logger.ErrorField(err), ) return "", nil, "FAILED", err.Error(), 5 * time.Second } } // --- Process Final Response --- if statusCode != http.StatusOK && statusCode != http.StatusCreated { errMsg = fmt.Sprintf("HTTP %d", statusCode) logger.Default().Error("[IMAGING STUDY WORKER] ❌ Respons Error dari API Part 3", 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 } // --- Success --- 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("[IMAGING STUDY WORKER] ✅ Sukses mengirim data Imaging Study", logger.String("source_id", sourceID), logger.String("fhir_id", fhirID), ) return fhirID, respBody, "SUCCESS", "", 0 } func (w *worker) appendImagingStudyMapping(createdAt time.Time, sourceID, fhirID, serviceRequestID, encounterID, patientID, practitionerID string) { // Format: YYYY-MM dirPath := filepath.Join("internal/worker/satusehat/imagingstudy/mapperdata", createdAt.Format("2006-01")) // Format: YYYY-MM-DD.txt fileName := createdAt.Format("2006-01-02") + ".txt" filePath := filepath.Join(dirPath, fileName) // Buat direktori jika belum ada if err := os.MkdirAll(dirPath, 0755); err != nil { logger.Default().Error("[IMAGING STUDY WORKER] Gagal membuat direktori mapping", logger.ErrorField(err), logger.String("path", dirPath)) return } // Buka file (buat jika belum ada, tambahkan jika sudah ada) f, err := os.OpenFile(filePath, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0644) if err != nil { logger.Default().Error("[IMAGING STUDY WORKER] Gagal membuka/membuat file mapping imagingstudy", logger.ErrorField(err), logger.String("file", filePath)) return } defer f.Close() // Tulis data ke file line := fmt.Sprintf("%s|%s|%s|%s|%s|%s\n", sourceID, fhirID, serviceRequestID, encounterID, patientID, practitionerID) if _, err := f.WriteString(line); err != nil { logger.Default().Error("[IMAGING STUDY WORKER] Gagal menulis ke file mapping", logger.ErrorField(err), logger.String("file", filePath)) } } func (w *worker) updateTracker(lastCreatedAt *time.Time, currentCreatedAt time.Time) { *lastCreatedAt = currentCreatedAt os.WriteFile(trackerFile, []byte(currentCreatedAt.Format(time.RFC3339Nano)), 0644) } func (w *worker) readLastCreatedAt() time.Time { data, err := os.ReadFile(trackerFile) if err != nil { // Jika file tidak ada, mulai dari awal (zero time) return time.Time{} } timestampStr := strings.TrimSpace(string(data)) t, err := time.Parse(time.RFC3339Nano, timestampStr) if err != nil { logger.Default().Warn("[IMAGING STUDY WORKER] Gagal parsing timestamp dari tracker file, memulai dari awal.", logger.String("content", timestampStr), logger.ErrorField(err)) return time.Time{} } return t } func truncate(s string, max int) string { if len(s) > max { return s[:max] + "..." } return s }