package encounter import ( "context" "database/sql" "encoding/json" "fmt" "log" "net/http" "os" "strconv" "strings" "time" "service/internal/infrastructure/database" extapi "service/internal/worker/interface" ) const trackerFile = "last_migrated_encounter_id.txt" // Config menyimpan konfigurasi untuk worker Encounter type Config struct { DBManager database.Service InternalBaseURL string // URL Internal API (contoh: http://localhost:8080/api/v1) InternalToken string // Token jika endpoint internal dilindungi Auth OrganizationID string } // WorkerService adalah interface untuk trigger scheduler type WorkerService interface { Run(ctx context.Context) } type worker struct { cfg Config repo Repository apiClient extapi.Client } // NewWorker inisiasi worker service baru func NewWorker(cfg Config) WorkerService { return &worker{ cfg: cfg, repo: NewRepository(cfg.DBManager), apiClient: extapi.NewClient(15 * time.Second), } } // Run adalah fungsi utama worker/scheduler func (w *worker) Run(ctx context.Context) { log.Println("[ENCOUNTER WORKER] Memulai proses sinkronisasi ke API Satu Sehat...") lastID := w.readLastID() log.Printf("[ENCOUNTER WORKER] Melanjutkan migrasi dari idxdaftar > %d\n", lastID) for { select { case <-ctx.Done(): log.Println("[ENCOUNTER WORKER] Proses dihentikan oleh sistem.") return default: } pendaftaran, err := w.repo.GetNextPendaftaran(ctx, lastID) if err != nil { if err == sql.ErrNoRows { time.Sleep(5 * time.Second) // Tunggu data baru jika habis continue } log.Printf("[ENCOUNTER WORKER] Gagal mengambil data database: %v\n", err) time.Sleep(5 * time.Second) continue } log.Printf("[ENCOUNTER WORKER] Memproses Pendaftaran ID: %d, RM: %s\n", pendaftaran.IdxDaftar, pendaftaran.NoMR.String) // Fallback pencarian IHS Number lewat NIK jika data master kosong if !pendaftaran.PatientIHS.Valid || pendaftaran.PatientIHS.String == "" { if pendaftaran.NIK.Valid && pendaftaran.NIK.String != "" { log.Printf("[ENCOUNTER WORKER] Patient IHS tidak ditemukan untuk ID %d. Mencoba mencari dari NIK: %s\n", pendaftaran.IdxDaftar, pendaftaran.NIK.String) searchURL := fmt.Sprintf("%s/satusehat/reference/patient/nik/%s", strings.TrimRight(w.cfg.InternalBaseURL, "/"), pendaftaran.NIK.String) searchResp, searchCode, searchErr := w.apiClient.GetJSON(ctx, searchURL, w.cfg.InternalToken) if searchErr == nil && searchCode == http.StatusOK { var searchResult map[string]interface{} if err := json.Unmarshal(searchResp, &searchResult); err == nil { // Ekstrak ID (IHS Number) dari atribut "data.id" pada format API standar Anda if data, ok := searchResult["data"].(map[string]interface{}); ok { if fhirID, ok := data["id"].(string); ok && fhirID != "" { pendaftaran.PatientIHS.String = fhirID pendaftaran.PatientIHS.Valid = true log.Printf("[ENCOUNTER WORKER] Berhasil mendapatkan IHS %s dari pencarian NIK.\n", fhirID) } } } } } } // Fallback pencarian Practitioner IHS Number lewat NIK jika data master kosong if !pendaftaran.PractitionerIHS.Valid || pendaftaran.PractitionerIHS.String == "" { if pendaftaran.PractitionerNIK.Valid && pendaftaran.PractitionerNIK.String != "" { log.Printf("[ENCOUNTER WORKER] Practitioner IHS tidak ditemukan untuk ID %d. Mencoba mencari dari NIK: %s\n", pendaftaran.IdxDaftar, pendaftaran.PractitionerNIK.String) searchURL := fmt.Sprintf("%s/satusehat/reference/practitioner/nik/%s", strings.TrimRight(w.cfg.InternalBaseURL, "/"), pendaftaran.PractitionerNIK.String) searchResp, searchCode, searchErr := w.apiClient.GetJSON(ctx, searchURL, w.cfg.InternalToken) if searchErr == nil && searchCode == http.StatusOK { var searchResult map[string]interface{} if err := json.Unmarshal(searchResp, &searchResult); err == nil { if data, ok := searchResult["data"].(map[string]interface{}); ok { if fhirID, ok := data["id"].(string); ok && fhirID != "" { pendaftaran.PractitionerIHS.String = fhirID pendaftaran.PractitionerIHS.Valid = true log.Printf("[ENCOUNTER WORKER] Berhasil mendapatkan Practitioner IHS %s dari pencarian NIK.\n", fhirID) } } } } } } // Skip record jika data IHS masih belum ada (gagal dari DB dan gagal dicari dari NIK) if !pendaftaran.PatientIHS.Valid || pendaftaran.PatientIHS.String == "" { log.Printf("[ENCOUNTER WORKER] Pendaftaran ID %d di-skip karena Patient IHS tidak ditemukan dan pencarian NIK gagal/kosong.\n", pendaftaran.IdxDaftar) w.updateTracker(&lastID, pendaftaran.IdxDaftar) continue } internalPayload := MapToInternalAPI(pendaftaran, w.cfg.OrganizationID) // Mengirim data ke endpoint internal API Faskes (POST /satusehat/encounter) reqURL := fmt.Sprintf("%s/satusehat/encounter", strings.TrimRight(w.cfg.InternalBaseURL, "/")) // Eksekusi HTTP request menggunakan Interface external client yang sudah dibuat respBody, statusCode, err := w.apiClient.PostJSON(ctx, reqURL, w.cfg.InternalToken, internalPayload) // Evaluasi Response dan Ekstrak ID (Jika Sukses) var status, errMsg, fhirID string if err != nil { status = "FAILED" errMsg = err.Error() log.Printf("[ENCOUNTER WORKER] Gagal HTTP Request ke Internal API (ID %d): %v\n", pendaftaran.IdxDaftar, err) } else if statusCode != http.StatusOK && statusCode != http.StatusCreated { status = "FAILED" errMsg = fmt.Sprintf("HTTP %d", statusCode) log.Printf("[ENCOUNTER WORKER] Gagal Response Internal API ID %d. HTTP: %d, Response: %s\n", pendaftaran.IdxDaftar, statusCode, string(respBody)) } else { status = "SUCCESS" // Ambil ID FHIR Kemenkes dari response format Standar Golang Faskes (Result.Data) 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 } } log.Printf("[ENCOUNTER WORKER] Sukses mengirim ID %d. Response: %s\n", pendaftaran.IdxDaftar, truncate(string(respBody), 150)) } // Simpan histori/log ke database untuk kemudahan audit dan resend/rollback syncLog := MapToSyncLog(pendaftaran.IdxDaftar, fhirID, internalPayload, respBody, status, errMsg) _ = w.repo.SaveSyncLog(ctx, syncLog) // Jika sukses, simpan ID FHIR ke tabel transaksi utama if status == "SUCCESS" && fhirID != "" { _ = w.repo.SaveSatuSehatID(ctx, pendaftaran.IdxDaftar, fhirID) } // Update file tracker (PENTING: Lanjut ke ID berikutnya agar tidak stuck pada data yang selalu error) w.updateTracker(&lastID, pendaftaran.IdxDaftar) // Jika gagal, tahan sebentar sebelum lanjut (Rate limiting/jeda darurat) if status == "FAILED" { time.Sleep(5 * time.Second) } else { time.Sleep(500 * time.Millisecond) } } } func (w *worker) updateTracker(lastID *int64, currentID int64) { *lastID = currentID w.saveLastID(*lastID) } // --- Tracker File Management --- func (w *worker) readLastID() int64 { data, err := os.ReadFile(trackerFile) if err != nil { return 0 } idStr := strings.TrimSpace(string(data)) id, err := strconv.ParseInt(idStr, 10, 64) if err != nil { return 0 } return id } func (w *worker) saveLastID(id int64) { os.WriteFile(trackerFile, []byte(strconv.FormatInt(id, 10)), 0644) } func truncate(s string, max int) string { if len(s) > max { return s[:max] + "..." } return s }