Files
satusehat-worker/internal/master/kfa/puller.go
T
2026-08-05 06:47:59 +00:00

219 lines
5.9 KiB
Go

package kfa
import (
"context"
"encoding/json"
"fmt"
"io"
"net/http"
"os"
"strconv"
"strings"
"time"
"service/pkg/logger"
)
type TokenManager interface {
GetAccessToken() string
ForceRefreshAndGetToken(ctx context.Context) (string, error)
}
type Config struct {
InternalBaseURL string
PageSize int
}
type Worker struct {
cfg Config
repo CommandRepository
tokenManager TokenManager
}
func NewWorker(cfg Config, repo CommandRepository, tokenManager TokenManager) *Worker {
if cfg.PageSize == 0 {
cfg.PageSize = 100
}
return &Worker{
cfg: cfg,
repo: repo,
tokenManager: tokenManager,
}
}
func (w *Worker) readTracker() int {
b, err := os.ReadFile("internal/master/kfa/kfa_tracker.txt")
if err != nil {
return 1
}
p, err := strconv.Atoi(strings.TrimSpace(string(b)))
if err != nil || p < 1 {
return 1
}
return p
}
func (w *Worker) writeTracker(page int) {
os.WriteFile("internal/master/kfa/kfa_tracker.txt", []byte(strconv.Itoa(page)), 0644)
}
func (w *Worker) Run(ctx context.Context) {
page := w.readTracker()
for {
select {
case <-ctx.Done():
logger.Default().Info("KFA Master Puller stopped")
return
default:
err := w.processPage(ctx, page)
if err != nil {
logger.Default().Error("KFA Master Puller error processing page", logger.ErrorField(err), logger.Int("page", page))
time.Sleep(30 * time.Second)
continue
}
page++
w.writeTracker(page)
time.Sleep(5 * time.Second) // Jedah halus antar halaman
}
}
}
func (w *Worker) processPage(ctx context.Context, page int) error {
listURL := fmt.Sprintf("%s/satusehat/reference/kfa/products?page=%d&size=%d&product_type=farmasi", strings.TrimRight(w.cfg.InternalBaseURL, "/"), page, w.cfg.PageSize)
req, err := http.NewRequestWithContext(ctx, http.MethodGet, listURL, nil)
if err != nil {
return err
}
req.Header.Set("Authorization", "Bearer "+w.tokenManager.GetAccessToken())
req.Header.Set("Accept", "application/json")
client := &http.Client{Timeout: 30 * time.Second}
resp, err := client.Do(req)
if err != nil {
return err
}
if resp.StatusCode == http.StatusUnauthorized {
resp.Body.Close()
logger.Default().Warn("KFA Master Puller list API 401, refreshing token...", logger.Int("page", page))
newToken, err := w.tokenManager.ForceRefreshAndGetToken(ctx)
if err != nil {
return err
}
req, _ = http.NewRequestWithContext(ctx, http.MethodGet, listURL, nil)
req.Header.Set("Authorization", "Bearer "+newToken)
req.Header.Set("Accept", "application/json")
resp, err = client.Do(req)
if err != nil {
return err
}
}
defer resp.Body.Close()
if resp.StatusCode == http.StatusTooManyRequests {
logger.Default().Warn("[KFA WORKER] ⏸️ Rate limit (429) pada List API, jeda 30s & ulangi halaman", logger.Int("page", page))
return fmt.Errorf("rate limited on list API") // Akan memicu waktu tidur 30 detik di method Run() dan mencoba halaman yang sama
}
if resp.StatusCode != http.StatusOK {
return fmt.Errorf("list API returned status: %d", resp.StatusCode)
}
body, err := io.ReadAll(resp.Body)
if err != nil {
return err
}
var listRes KfaListResponse
if err := json.Unmarshal(body, &listRes); err != nil {
return err
}
if len(listRes.Data.Items.Data) == 0 {
logger.Default().Info("No more KFA products found, resetting to page 1 and sleeping", logger.Int("page", page))
w.writeTracker(1)
time.Sleep(6 * time.Hour)
return fmt.Errorf("empty page, reset triggered")
}
logger.Default().Info("Pulling KFA items", logger.Int("page", page), logger.Int("items", len(listRes.Data.Items.Data)))
for _, item := range listRes.Data.Items.Data {
for {
rateLimited, err := w.fetchAndSaveDetail(ctx, item.KfaCode)
if rateLimited {
logger.Default().Warn("[KFA WORKER] ⏸️ Rate limit (429) dari Satu Sehat - jeda 60 detik dan retry dokumen yang sama", logger.String("kfa_code", item.KfaCode))
time.Sleep(60 * time.Second)
continue // Ulangi request untuk item ini (Tidak melewatkan ID)
}
if err != nil {
logger.Default().Error("Failed to fetch/save KFA detail", logger.String("kfa_code", item.KfaCode), logger.ErrorField(err))
}
break // Lanjut ke item/ID berikutnya bila sukses atau gagal karena alasan non-ratelimit
}
time.Sleep(500 * time.Millisecond) // Pencegahan Rate Limit Part 3
}
return nil
}
func (w *Worker) fetchAndSaveDetail(ctx context.Context, kfaCode string) (bool, error) {
detailURL := fmt.Sprintf("%s/satusehat/reference/kfa/products/%s", strings.TrimRight(w.cfg.InternalBaseURL, "/"), kfaCode)
req, err := http.NewRequestWithContext(ctx, http.MethodGet, detailURL, nil)
if err != nil {
return false, err
}
req.Header.Set("Authorization", "Bearer "+w.tokenManager.GetAccessToken())
req.Header.Set("Accept", "application/json")
client := &http.Client{Timeout: 30 * time.Second}
resp, err := client.Do(req)
if err != nil {
return false, err
}
if resp.StatusCode == http.StatusUnauthorized {
resp.Body.Close()
newToken, err := w.tokenManager.ForceRefreshAndGetToken(ctx)
if err != nil {
return false, err
}
req, _ = http.NewRequestWithContext(ctx, http.MethodGet, detailURL, nil)
req.Header.Set("Authorization", "Bearer "+newToken)
req.Header.Set("Accept", "application/json")
resp, err = client.Do(req)
if err != nil {
return false, err
}
}
defer resp.Body.Close()
if resp.StatusCode == http.StatusTooManyRequests {
return true, fmt.Errorf("rate limited (429)")
}
if resp.StatusCode != http.StatusOK {
return false, fmt.Errorf("detail API returned status: %d", resp.StatusCode)
}
body, err := io.ReadAll(resp.Body)
if err != nil {
return false, err
}
var detailRes KfaDetailResponse
if err := json.Unmarshal(body, &detailRes); err != nil {
return false, err
}
product, ingredients, packagings := MapKfaDetailToEntity(detailRes.Data.Result)
return false, w.repo.UpsertProductBundle(ctx, product, ingredients, packagings)
}