package imagingstudy import ( "context" "database/sql" "encoding/json" "fmt" "net/http" "os" "strings" "time" "service/internal/infrastructure/database" extapi "service/internal/worker/interface" "service/pkg/logger" ) const trackerFile = "internal/worker/satusehat/imagingstudy/last_imagingstudy_id.txt" 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, shouldRetry := w.sendToInternalAPI(ctx, internalPayload, imagingData.SourceID) /* if err != nil { // Kegagalan Jaringan/Timeout status = "FAILED" errMsg = err.Error() logger.Default().Error("[IMAGING STUDY WORKER] ❌ Gagal Koneksi ke API Part 3", logger.String("source_id", imagingData.SourceID), logger.String("no_mr", imagingData.NoMR.String), logger.Any("payload_sent", internalPayload), logger.ErrorField(err), ) shouldRetry = true // Wajib retry jika jaringan terputus/timeout } else if statusCode != http.StatusOK && statusCode != http.StatusCreated { // Kegagalan Response API status = "FAILED" errMsg = fmt.Sprintf("HTTP %d", statusCode) logger.Default().Error("[IMAGING STUDY WORKER] ❌ Respons Error dari API Part 3", logger.String("source_id", imagingData.SourceID), logger.Int("status_code", statusCode), logger.String("response", truncate(string(respBody), 200)), ) // Retry jika Internal Server Error (5xx) atau Rate Limit (429) if statusCode >= 500 || statusCode == http.StatusTooManyRequests { shouldRetry = true } } else { // Sukses status = "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", imagingData.SourceID), logger.String("fhir_id", fhirID), ) }*/ // 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)) } } // Update Tracker TXT agar tidak dikirim ulang terus-menerus KECUALI jika butuh retry if !shouldRetry { w.updateTracker(&lastCreatedAt, imagingData.CreatedAt) } else { // Memberitahu log bahwa data yang sama akan dicoba lagi pada loop berikutnya logger.Default().Warn("[IMAGING STUDY WORKER] ⚠️ Retrying data", logger.String("source_id", imagingData.SourceID), logger.Time("created_at", imagingData.CreatedAt), logger.String("reason", errMsg)) } // Jeda if shouldRetry { time.Sleep(5 * time.Second) // Tunggu sebentar jika server tujuan sedang sibuk } else { time.Sleep(2 * time.Second) // Jeda antar pengiriman diperpanjang untuk mencegah Rate Limit } } } // 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, shouldRetry bool) { 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(), true } // --- 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(), true } 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(), true } } // --- 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)), ) // Retry on server-side errors if statusCode >= 500 || statusCode == http.StatusTooManyRequests { return "", respBody, "FAILED", errMsg, true } // Do not retry on client-side errors (4xx) return "", respBody, "FAILED", errMsg, false } // --- 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", "", false } 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 }