first commit
This commit is contained in:
@@ -0,0 +1,28 @@
|
||||
package imagingstudy
|
||||
|
||||
import (
|
||||
"database/sql"
|
||||
)
|
||||
|
||||
// ImagingStudyData merepresentasikan gabungan data dari Postgres (Satu Data) dan MySQL (RIS)
|
||||
type ImagingStudyData struct {
|
||||
RequestID string `db:"id"`
|
||||
ServiceRequestID sql.NullString `db:"Entity_Id"`
|
||||
EncounterID sql.NullString `db:"Encounter_Id"`
|
||||
NoRegister sql.NullString `db:"No_register"`
|
||||
PatientID sql.NullString `db:"Nomor_satusehat_pasien"`
|
||||
PatientName sql.NullString `db:"Nama_lengkap"`
|
||||
|
||||
// Data dari MySQL (RIS)
|
||||
Modality sql.NullString
|
||||
StartedDate sql.NullTime
|
||||
}
|
||||
|
||||
type ImagingStudySyncLog struct {
|
||||
RequestID string `db:"id"`
|
||||
ImagingStudyID string `db:"imagingstudy_id"`
|
||||
RequestPayload string `db:"request_payload"`
|
||||
ResponsePayload string `db:"response_payload"`
|
||||
Status string `db:"status"`
|
||||
ErrorMessage string `db:"error_message"`
|
||||
}
|
||||
@@ -0,0 +1,36 @@
|
||||
package imagingstudy
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"time"
|
||||
)
|
||||
|
||||
// MapToInternalAPI memetakan data gabungan ke payload API Internal SatuSehat
|
||||
func MapToInternalAPI(dbData *ImagingStudyData, orgID string) map[string]interface{} {
|
||||
waktu := time.Now()
|
||||
if dbData.StartedDate.Valid {
|
||||
waktu = dbData.StartedDate.Time
|
||||
}
|
||||
payload := map[string]interface{}{
|
||||
"request_id": dbData.RequestID,
|
||||
"servicerequest_id": dbData.ServiceRequestID.String,
|
||||
"encounter_id": dbData.EncounterID.String,
|
||||
"patient_id": dbData.PatientID.String,
|
||||
"accession_number": dbData.NoRegister.String, // Digunakan sebagai Identifier SATUSEHAT
|
||||
"modality": dbData.Modality.String,
|
||||
"started_date": waktu.Format(time.RFC3339),
|
||||
}
|
||||
return payload
|
||||
}
|
||||
|
||||
func MapToSyncLog(reqID string, fhirID string, reqPayload interface{}, respBody []byte, status string, errMsg string) ImagingStudySyncLog {
|
||||
reqBytes, _ := json.MarshalIndent(reqPayload, "", " ")
|
||||
return ImagingStudySyncLog{
|
||||
RequestID: reqID,
|
||||
ImagingStudyID: fhirID,
|
||||
RequestPayload: string(reqBytes),
|
||||
ResponsePayload: string(respBody),
|
||||
Status: status,
|
||||
ErrorMessage: errMsg,
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,87 @@
|
||||
package imagingstudy
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
|
||||
"service/internal/infrastructure/database"
|
||||
)
|
||||
|
||||
type Repository interface {
|
||||
GetNextPendingRequest(ctx context.Context) (*ImagingStudyData, error)
|
||||
UpdateSyncStatus(ctx context.Context, reqID string, status string, responsePayload string) error
|
||||
}
|
||||
|
||||
type repository struct {
|
||||
dbManager database.Service
|
||||
}
|
||||
|
||||
func NewRepository(dbManager database.Service) Repository {
|
||||
return &repository{dbManager: dbManager}
|
||||
}
|
||||
|
||||
// GetNextPendingRequest mengambil 1 data service request radiologi yang belum terkirim
|
||||
func (r *repository) GetNextPendingRequest(ctx context.Context) (*ImagingStudyData, error) {
|
||||
var data ImagingStudyData
|
||||
|
||||
// 1. Ambil koneksi Postgres (Satu Data) - Sesuaikan nama dengan konfigurasi Anda
|
||||
pgDB, err := r.dbManager.GetDB("satudata")
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// Ambil dari Postgres
|
||||
queryPG := `
|
||||
SELECT
|
||||
sr.id, sr."Entity_Id", sr."Encounter_Id", sr."No_register",
|
||||
kp."Nomor_satusehat_pasien", kp."Nama_lengkap"
|
||||
FROM public.data_service_requests sr
|
||||
JOIN public.data_kunjungan_pasien kp ON sr."Encounter_Id" = kp."IDXDAFTAR_satusehat"
|
||||
WHERE LOWER(sr."Order_from") = 'radiologi'
|
||||
AND (sr."Send_status" IS NULL OR sr."Send_status" = 'FAILED')
|
||||
ORDER BY sr."Created_at" ASC
|
||||
LIMIT 1
|
||||
`
|
||||
err = pgDB.QueryRowContext(ctx, queryPG).Scan(
|
||||
&data.RequestID, &data.ServiceRequestID, &data.EncounterID, &data.NoRegister,
|
||||
&data.PatientID, &data.PatientName,
|
||||
)
|
||||
if err != nil { // Bisa berupa sql.ErrNoRows jika data habis
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// 2. Ambil koneksi MySQL (RIS) - Sesuaikan nama dengan konfigurasi Anda
|
||||
mySQLDB, err := r.dbManager.GetDB("simrs")
|
||||
if err == nil && data.NoRegister.Valid {
|
||||
queryMySQL := `SELECT modality, foto FROM periksa WHERE noregister = ? LIMIT 1`
|
||||
|
||||
var modality sql.NullString
|
||||
var foto sql.NullTime
|
||||
|
||||
// Gunakan placeholder '?' jika MySQL, atau '$1' jika PostgreSQL
|
||||
err = mySQLDB.QueryRowContext(ctx, queryMySQL, data.NoRegister.String).Scan(&modality, &foto)
|
||||
if err == nil {
|
||||
data.Modality = modality
|
||||
data.StartedDate = foto
|
||||
}
|
||||
}
|
||||
|
||||
return &data, nil
|
||||
}
|
||||
|
||||
// UpdateSyncStatus memperbarui status pengiriman di PostgreSQL
|
||||
func (r *repository) UpdateSyncStatus(ctx context.Context, reqID string, status string, responsePayload string) error {
|
||||
pgDB, err := r.dbManager.GetDB("satudata")
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
query := `
|
||||
UPDATE public.data_service_requests
|
||||
SET "Send_status" = $1, "Send_response" = $2, "Updated_at" = NOW()
|
||||
WHERE id = $3
|
||||
`
|
||||
|
||||
_, err = pgDB.ExecContext(ctx, query, status, responsePayload, reqID)
|
||||
return err
|
||||
}
|
||||
@@ -0,0 +1,98 @@
|
||||
package imagingstudy
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"log"
|
||||
"net/http"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"service/internal/infrastructure/database"
|
||||
extapi "service/internal/worker/interface"
|
||||
)
|
||||
|
||||
type Config struct {
|
||||
DBManager database.Service
|
||||
InternalBaseURL string
|
||||
InternalToken string
|
||||
OrganizationID string
|
||||
}
|
||||
|
||||
type WorkerService interface {
|
||||
Run(ctx context.Context)
|
||||
}
|
||||
|
||||
type worker struct {
|
||||
cfg Config
|
||||
repo Repository
|
||||
apiClient extapi.Client
|
||||
}
|
||||
|
||||
func NewWorker(cfg Config) WorkerService {
|
||||
return &worker{cfg: cfg, repo: NewRepository(cfg.DBManager), apiClient: extapi.NewClient(15 * time.Second)}
|
||||
}
|
||||
|
||||
func (w *worker) Run(ctx context.Context) {
|
||||
log.Println("[IMAGINGSTUDY WORKER] Memulai proses sinkronisasi...")
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
log.Println("[IMAGINGSTUDY WORKER] Proses dihentikan.")
|
||||
return
|
||||
default:
|
||||
}
|
||||
|
||||
// 1. Ambil data dengan status belum terkirim (IS NULL atau FAILED)
|
||||
isData, err := w.repo.GetNextPendingRequest(ctx)
|
||||
if err != nil {
|
||||
if err == sql.ErrNoRows {
|
||||
// Jika tidak ada data antrean, istirahat lebih lama
|
||||
time.Sleep(10 * time.Second)
|
||||
continue
|
||||
}
|
||||
log.Printf("[IMAGINGSTUDY WORKER] Error query db: %v", err)
|
||||
time.Sleep(5 * time.Second)
|
||||
continue
|
||||
}
|
||||
|
||||
// 2. Petakan ke payload SATUSEHAT via Internal API
|
||||
internalPayload := MapToInternalAPI(isData, w.cfg.OrganizationID)
|
||||
reqURL := fmt.Sprintf("%s/satusehat/imagingstudy", strings.TrimRight(w.cfg.InternalBaseURL, "/"))
|
||||
|
||||
// 3. Eksekusi pengiriman API
|
||||
respBody, statusCode, err := w.apiClient.PostJSON(ctx, reqURL, w.cfg.InternalToken, internalPayload)
|
||||
|
||||
var status, errMsg, fhirID string
|
||||
if err != nil || (statusCode != http.StatusOK && statusCode != http.StatusCreated) {
|
||||
status, errMsg = "FAILED", fmt.Sprintf("Err: %v, HTTP: %d", err, statusCode)
|
||||
} else {
|
||||
status = "SUCCESS"
|
||||
var res map[string]interface{}
|
||||
if json.Unmarshal(respBody, &res) == nil {
|
||||
if data, ok := res["data"].(map[string]interface{}); ok {
|
||||
fhirID, _ = data["id"].(string)
|
||||
} else {
|
||||
fhirID, _ = res["id"].(string)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Buat response data berupa string JSON untuk log & penyimpanan error
|
||||
logData := MapToSyncLog(isData.RequestID, fhirID, internalPayload, respBody, status, errMsg)
|
||||
logBytes, _ := json.Marshal(logData)
|
||||
|
||||
// 4. Update status ke tabel data_service_requests
|
||||
errUpdate := w.repo.UpdateSyncStatus(ctx, isData.RequestID, status, string(logBytes))
|
||||
if errUpdate != nil {
|
||||
log.Printf("[IMAGINGSTUDY WORKER] Gagal update status di database untuk ID %s: %v", isData.RequestID, errUpdate)
|
||||
} else {
|
||||
log.Printf("[IMAGINGSTUDY WORKER] Sukses proses ID %s dengan status %s", isData.RequestID, status)
|
||||
}
|
||||
|
||||
time.Sleep(500 * time.Millisecond)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user