first commit
This commit is contained in:
@@ -0,0 +1,28 @@
|
||||
package observation
|
||||
|
||||
import (
|
||||
"database/sql"
|
||||
)
|
||||
|
||||
// ObservationDB merepresentasikan struktur tabel pemeriksaan (Contoh: Vital Sign/Lab)
|
||||
// TODO: Sesuaikan field dan tag `db` dengan tabel SIMRS Anda yang sebenarnya
|
||||
type ObservationDB struct {
|
||||
IdxPemeriksaan int64 `db:"idxpemeriksaan"` // Primary Key
|
||||
IdxDaftar int64 `db:"idxdaftar"` // Relasi ke Encounter
|
||||
NoMR sql.NullString `db:"nomr"`
|
||||
TglPeriksa sql.NullTime `db:"tgl_periksa"`
|
||||
KodePemeriksaan sql.NullString `db:"kode_pemeriksaan"` // Contoh: LOINC Code
|
||||
Nilai sql.NullString `db:"nilai"`
|
||||
Satuan sql.NullString `db:"satuan"`
|
||||
KdDokter sql.NullInt64 `db:"kddokter"`
|
||||
}
|
||||
|
||||
// ObservationSyncLog merepresentasikan struktur log sinkronisasi API SatuSehat.
|
||||
type ObservationSyncLog struct {
|
||||
IdxPemeriksaan int64 `db:"idxpemeriksaan"`
|
||||
ObservationID string `db:"observation_id"` // ID dari SatuSehat jika SUCCESS
|
||||
RequestPayload string `db:"request_payload"` // Teks JSON raw request
|
||||
ResponsePayload string `db:"response_payload"` // Teks JSON raw response
|
||||
Status string `db:"status"` // Status: SUCCESS / FAILED
|
||||
ErrorMessage string `db:"error_message"` // Pesan kegagalan
|
||||
}
|
||||
@@ -0,0 +1,42 @@
|
||||
package observation
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"time"
|
||||
)
|
||||
|
||||
// MapToInternalAPI memetakan data database menjadi payload untuk Internal API GoPrint
|
||||
// TODO: Sesuaikan dengan JSON Payload yang dibutuhkan oleh internal endpoint Observation Anda
|
||||
func MapToInternalAPI(dbData *ObservationDB, orgID string) map[string]interface{} {
|
||||
// Format Waktu
|
||||
waktuPeriksa := time.Now()
|
||||
if dbData.TglPeriksa.Valid {
|
||||
waktuPeriksa = dbData.TglPeriksa.Time
|
||||
}
|
||||
|
||||
payload := map[string]interface{}{
|
||||
"observation_id": fmt.Sprintf("%d", dbData.IdxPemeriksaan),
|
||||
"encounter_id": fmt.Sprintf("%d", dbData.IdxDaftar),
|
||||
"organization_id": orgID,
|
||||
"patient_id": dbData.NoMR.String, // Sebaiknya diganti dengan IHS Number di Service
|
||||
"code": dbData.KodePemeriksaan.String,
|
||||
"value": dbData.Nilai.String,
|
||||
"unit": dbData.Satuan.String,
|
||||
"effective_date": waktuPeriksa.Format(time.RFC3339),
|
||||
}
|
||||
return payload
|
||||
}
|
||||
|
||||
// MapToSyncLog membuat struktur data log histori untuk disimpan ke tabel database
|
||||
func MapToSyncLog(idx int64, fhirID string, reqPayload interface{}, respBody []byte, status string, errMsg string) ObservationSyncLog {
|
||||
reqBytes, _ := json.MarshalIndent(reqPayload, "", " ")
|
||||
return ObservationSyncLog{
|
||||
IdxPemeriksaan: idx,
|
||||
ObservationID: fhirID,
|
||||
RequestPayload: string(reqBytes),
|
||||
ResponsePayload: string(respBody),
|
||||
Status: status,
|
||||
ErrorMessage: errMsg,
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,92 @@
|
||||
package observation
|
||||
|
||||
import (
|
||||
"context"
|
||||
"service/internal/infrastructure/database"
|
||||
"service/pkg/utils/query"
|
||||
)
|
||||
|
||||
type Repository interface {
|
||||
GetNextObservation(ctx context.Context, lastID int64) (*ObservationDB, error)
|
||||
SaveSatuSehatID(ctx context.Context, idx int64, fhirID string) error
|
||||
SaveSyncLog(ctx context.Context, logData ObservationSyncLog) error
|
||||
}
|
||||
|
||||
type repository struct {
|
||||
dbManager database.Service
|
||||
}
|
||||
|
||||
func NewRepository(dbManager database.Service) Repository {
|
||||
return &repository{dbManager: dbManager}
|
||||
}
|
||||
|
||||
func (r *repository) GetNextObservation(ctx context.Context, lastID int64) (*ObservationDB, error) {
|
||||
var obs ObservationDB
|
||||
simrsDB, err := r.dbManager.GetSQLXDB("simrs")
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
qb := query.NewSQLQueryBuilder(query.DBTypePostgreSQL).SetSecurityOptions(false, 0)
|
||||
|
||||
// TODO: Sesuaikan dengan nama tabel pemeriksaan di SIMRS Anda
|
||||
qSimrs := query.DynamicQuery{
|
||||
From: "public.t_pemeriksaan",
|
||||
Filters: []query.FilterGroup{
|
||||
{
|
||||
Filters: []query.DynamicFilter{
|
||||
query.CreateFilter("idxpemeriksaan", query.OpGreaterThan, lastID),
|
||||
},
|
||||
},
|
||||
},
|
||||
Sort: []query.SortField{
|
||||
query.CreateAscSort("idxpemeriksaan"),
|
||||
},
|
||||
Limit: 1,
|
||||
}
|
||||
|
||||
err = qb.ExecuteQueryRow(ctx, simrsDB, qSimrs, &obs)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &obs, nil
|
||||
}
|
||||
|
||||
func (r *repository) SaveSatuSehatID(ctx context.Context, idx int64, fhirID string) error {
|
||||
simrsDB, err := r.dbManager.GetSQLXDB("simrs")
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
qb := query.NewSQLQueryBuilder(query.DBTypePostgreSQL).SetSecurityOptions(false, 0)
|
||||
updateData := query.UpdateData{
|
||||
Columns: []string{"status_bridging"},
|
||||
Values: []interface{}{fhirID},
|
||||
}
|
||||
filters := []query.FilterGroup{
|
||||
query.CreateAndFilterGroup([]query.DynamicFilter{query.CreateEqualFilter("idxpemeriksaan", idx)}),
|
||||
}
|
||||
|
||||
_, err = qb.ExecuteUpdate(ctx, simrsDB, "public.t_pemeriksaan", updateData, filters)
|
||||
return err
|
||||
}
|
||||
|
||||
func (r *repository) SaveSyncLog(ctx context.Context, logData ObservationSyncLog) error {
|
||||
simrsDB, err := r.dbManager.GetSQLXDB("simrs")
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
qb := query.NewSQLQueryBuilder(query.DBTypePostgreSQL).SetSecurityOptions(false, 0)
|
||||
insertData := query.InsertData{
|
||||
Columns: []string{"idxpemeriksaan", "observation_id", "request_payload", "response_payload", "status", "error_message"},
|
||||
Values: []interface{}{
|
||||
logData.IdxPemeriksaan, logData.ObservationID, logData.RequestPayload,
|
||||
logData.ResponsePayload, logData.Status, logData.ErrorMessage,
|
||||
},
|
||||
}
|
||||
conflictCols := []string{"idxpemeriksaan"}
|
||||
updateCols := []string{"observation_id", "request_payload", "response_payload", "status", "error_message"}
|
||||
_, err = qb.ExecuteUpsert(ctx, simrsDB, "public.log_satusehat_observation", insertData, conflictCols, updateCols)
|
||||
return err
|
||||
}
|
||||
@@ -0,0 +1,117 @@
|
||||
package observation
|
||||
|
||||
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_observation_id.txt"
|
||||
|
||||
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("[OBSERVATION WORKER] Memulai proses sinkronisasi ke API Satu Sehat...")
|
||||
lastID := w.readLastID()
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
log.Println("[OBSERVATION WORKER] Proses dihentikan oleh sistem.")
|
||||
return
|
||||
default:
|
||||
}
|
||||
|
||||
obs, err := w.repo.GetNextObservation(ctx, lastID)
|
||||
if err != nil {
|
||||
if err == sql.ErrNoRows {
|
||||
time.Sleep(5 * time.Second)
|
||||
continue
|
||||
}
|
||||
log.Printf("[OBSERVATION WORKER] Gagal mengambil data database: %v\n", err)
|
||||
time.Sleep(5 * time.Second)
|
||||
continue
|
||||
}
|
||||
|
||||
log.Printf("[OBSERVATION WORKER] Memproses Observasi ID: %d\n", obs.IdxPemeriksaan)
|
||||
internalPayload := MapToInternalAPI(obs, w.cfg.OrganizationID)
|
||||
reqURL := fmt.Sprintf("%s/satusehat/observation", strings.TrimRight(w.cfg.InternalBaseURL, "/"))
|
||||
|
||||
respBody, statusCode, err := w.apiClient.PostJSON(ctx, reqURL, w.cfg.InternalToken, internalPayload)
|
||||
|
||||
var status, errMsg, fhirID string
|
||||
if err != nil {
|
||||
status, errMsg = "FAILED", err.Error()
|
||||
log.Printf("[OBSERVATION WORKER] Gagal HTTP Request (ID %d): %v\n", obs.IdxPemeriksaan, err)
|
||||
} else if statusCode != http.StatusOK && statusCode != http.StatusCreated {
|
||||
status, errMsg = "FAILED", fmt.Sprintf("HTTP %d", statusCode)
|
||||
log.Printf("[OBSERVATION WORKER] Gagal Response ID %d. HTTP: %d, Body: %s\n", obs.IdxPemeriksaan, statusCode, string(respBody))
|
||||
} else {
|
||||
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
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
syncLog := MapToSyncLog(obs.IdxPemeriksaan, fhirID, internalPayload, respBody, status, errMsg)
|
||||
_ = w.repo.SaveSyncLog(ctx, syncLog)
|
||||
|
||||
if status == "SUCCESS" && fhirID != "" {
|
||||
_ = w.repo.SaveSatuSehatID(ctx, obs.IdxPemeriksaan, fhirID)
|
||||
}
|
||||
|
||||
w.updateTracker(&lastID, obs.IdxPemeriksaan)
|
||||
time.Sleep(500 * time.Millisecond) // Rate limiting
|
||||
}
|
||||
}
|
||||
|
||||
func (w *worker) updateTracker(lastID *int64, currentID int64) {
|
||||
*lastID = currentID
|
||||
os.WriteFile(trackerFile, []byte(strconv.FormatInt(*lastID, 10)), 0644)
|
||||
}
|
||||
|
||||
func (w *worker) readLastID() int64 {
|
||||
data, _ := os.ReadFile(trackerFile)
|
||||
id, _ := strconv.ParseInt(strings.TrimSpace(string(data)), 10, 64)
|
||||
return id
|
||||
}
|
||||
Reference in New Issue
Block a user