07. main.go: merakit semuanya
Bagianmu untuk bagian P4 dan polling. File asli: AAT/services/bnpb-aggregator/main.go.
Syntax: belajar_go/06_context_defer.go dan 07_goroutine_channel_select.go.
Kenapa ada
store, api, dan ingest adalah komponen-komponen terpisah. Harus ada satu tempat yang
membuat semuanya, menghubungkannya, dan menjalankannya. Itulah main.go. Ia tidak
berisi logika bisnis; logikanya ada di package internal/*.
Kenapa dibutuhkan
Tanpa perakitan, ketiga package itu ada tetapi tidak ada yang memanggilnya, dan service hanya
menjawab /health.
Gambaran alurnya
Mongo <-- store.Save <-- handle <-- ingest.Run <-- klien BMKG, klien PVMBG <-- mock
^
+--- store.List / Get <-- api (GET /hazards, /hazards/{id}) <-- bnpb-api
Bagian 1: main, pintu masuk
Dikutip dari main.go, baris 24 sampai 38
func main() {
core.RunHealthcheckIfRequested()
cfg, err := config.Load()
if err != nil {
fmt.Fprintf(os.Stderr, "invalid configuration:\n%v\n", err)
os.Exit(1)
}
logger := core.NewLogger(cfg.ServiceName, cfg.LogLevel)
if err := run(cfg, logger); err != nil {
logger.Error("service failed", "error", err.Error())
os.Exit(1)
}
}core.RunHealthcheckIfRequested(): kalau program dijalankan sebagai/app healthcheck(olehHEALTHCHECKDocker), ia memeriksa dirinya sendiri lalu keluar. Kalau tidak, tidak berefek. (Lihat03_pondasi_core.md.)config.Load(): baca konfigurasi. Kalau ada variabel wajib yang kosong, semua masalah dicetak sekaligus lalu keluar dengan kode 1.run(cfg, logger): seluruh pekerjaan ada di sini.
Kenapa dipisah menjadi main dan run? os.Exit() tidak menjalankan defer. Kalau
semuanya ada di main dan memanggil os.Exit saat gagal, koneksi Mongo tidak akan ditutup
rapi. Dengan run yang mengembalikan error, semua defer di dalamnya selesai dulu, baru
main memutuskan keluar. (Pelajaran 6.)
Bagian 2: run, urutan perakitan
Dikutip dari main.go, baris 40 sampai 102
func run(cfg config.Config, logger *slog.Logger) error {
// Cancelled on SIGTERM (docker compose stop) / Ctrl+C; background loops stop on it.
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGTERM, os.Interrupt)
defer stop()
// Canonical Store (P4). Only this service talks to Mongo.
connectCtx, cancel := context.WithTimeout(ctx, 15*time.Second)
st, err := store.New(connectCtx, cfg.MongoURI, cfg.MongoDB)
cancel()
if err != nil {
return fmt.Errorf("canonical store: %w", err)
}
defer func() {
closeCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
_ = st.Close(closeCtx)
}()
mux := core.NewMux(cfg.ServiceName)
api.Register(mux, logger, st) // internal read API, the only way other services read the store
// Polling. The source clients (P1/P3 branch) fetch from each agency and map the raw data
// to HazardEvents with the P1 mapper; ingest.Run only loops over them.
volcanoes, err := reference.Volcanoes()
if err != nil {
return err
}
mapper, err := mapping.NewMapper(volcanoes)
if err != nil {
return err
}
httpClient := core.NewHTTPClient(logger, cfg.UpstreamTimeout)
bmkgClient, err := bmkg.NewClient(cfg.BMKGBaseURL, cfg.BMKGAPIKey, httpClient, mapper)
if err != nil {
return err
}
pvmbgClient, err := pvmbg.NewClient(cfg.PVMBGBaseURL, cfg.PVMBGToken, httpClient, mapper)
if err != nil {
return err
}
clients := []ingest.SourceClient{bmkgClient, pvmbgClient}
// TODO(P5): also publish each event to Kafka here, after it is saved.
handle := func(ctx context.Context, events []domain.HazardEvent) error {
for _, e := range events {
if err := st.Save(ctx, e); err != nil {
return err
}
}
return nil
}
pollingDone := make(chan struct{})
go func() {
defer close(pollingDone)
ingest.Run(ctx, logger, cfg.PollInterval, clients, handle)
}()
serveErr := core.Serve(ctx, logger, cfg.Port, mux)
stop()
<-pollingDone // let the pollers finish their cycle before the store closes
return serveErr
}Urutannya, per langkah:
Langkah 1: ctx yang batal saat ada sinyal berhenti.
signal.NotifyContext(..., SIGTERM, os.Interrupt) membuat ctx yang otomatis dibatalkan
ketika Docker mengirim sinyal berhenti (docker compose stop) atau kamu menekan Ctrl+C. Semua
pekerjaan latar belakang memantau ctx ini lalu berhenti. Inilah dasar kriteria P4 poin 1
(satu container bisa dihentikan dengan rapi).
Langkah 2: sambung ke Mongo (store.New).
Dengan batas waktu 15 detik. Kalau gagal, run mengembalikan error: service tidak jalan
setengah-setengah, container keluar, dan Docker menjalankannya lagi (restart: unless-stopped).
Lebih baik gagal jelas di awal daripada berjalan tanpa penyimpanan.
Langkah 3: defer untuk menutup koneksi.
Koneksi Mongo ditutup saat run selesai, apa pun alasannya. Perhatikan penutupan memakai
context.Background() baru, bukan ctx di atas, karena ctx itu sudah dibatalkan saat
shutdown, dan penutupan koneksi masih butuh context yang hidup.
Langkah 4: pasang API (api.Register(mux, logger, st)).
st adalah *store.Store, tetapi Register menerima interface Reader. Karena *store.Store
punya method List dan Get, ia otomatis cocok. Setelah baris ini, GET /hazards dan
GET /hazards/{id} aktif.
Langkah 5: siapkan polling.
reference.Volcanoes(): tabel gunung api statis milik BNPB (di-embed di dalam binary).mapping.NewMapper(volcanoes): mapper milik Grace.core.NewHTTPClient(logger, cfg.UpstreamTimeout): HTTP client dengan timeout dan pencatatan latensi.bmkg.NewClient(...)danpvmbg.NewClient(...): klien milik Grace. Konstruktornya memeriksa masukan dan mengembalikan error, jadi salah konfigurasi ketahuan saat start. Masing-masing memakai kredensial sendiri.clients := []ingest.SourceClient{bmkgClient, pvmbgClient}: daftar klien (tipenya interface).
Langkah 6: handle, apa yang dilakukan dengan event.
handle := func(ctx context.Context, events []domain.HazardEvent) error {
for _, e := range events { if err := st.Save(ctx, e); err != nil { return err } }
return nil
}
Ini sebuah closure: fungsi tanpa nama yang mengingat st. Ini satu-satunya tempat yang
tahu bahwa event disimpan ke store. Kalau Save gagal, error dikembalikan, dan ingest tidak
menggeser cursor (retry). Komentar TODO(P5) menandai tempat Mike nanti mem-publish ke Kafka.
Langkah 7: jalankan polling di goroutine.
go func() { defer close(pollingDone); ingest.Run(...) }() memulai polling bersamaan dengan
server HTTP di bawahnya. pollingDone adalah channel yang ditutup saat ingest.Run selesai,
sebagai tanda.
Langkah 8: jalankan server dan matikan dengan rapi.
core.Serve(...) memblokir sampai ctx dibatalkan. Setelah itu:
stop()memastikanctxbatal, juga kalauServeberhenti karena error.<-pollingDonemenunggu poller benar-benar berhenti.return serveErr, laludeferdijalankan: Mongo ditutup.
Kenapa menunggu poller sebelum Mongo ditutup? Kalau Mongo ditutup duluan, goroutine poller bisa masih memakainya dan gagal di tengah jalan.
Urutan kejadian saat docker compose stop bnpb-aggregator
- Docker mengirim
SIGTERM. ctxdibatalkan.Serveberhenti menerima koneksi baru, memberi request yang berjalan sampai 10 detik.ingest.Runberhenti: panggilan HTTP yang berjalan dibatalkan, loop memantauctx.<-pollingDoneterbuka.defermenutup koneksi Mongo.- Proses keluar.
Setiap document disimpan secara atomik (satu upsert per document), jadi tidak ada document yang rusak. Satu batch bisa terpotong kalau siklus sedang menyimpan saat dihentikan; cursor tidak bergeser, jadi event itu diambil lagi dan di-upsert saat service menyala (tidak ganda).
Cara menjelaskan main.go dalam 30 detik
main.gomerakit service: sambung ke Mongo, pasang API baca, bangun dua klien sumber milik Grace, lalu jalankan polling di goroutine dengan handler yang menyimpan ke store. Semuanya berhenti dengan rapi saat ada sinyal: server berhenti, poller berhenti, dan koneksi Mongo ditutup paling akhir.
Latihan
- Kenapa
store.Newdipanggil sebelumapi.Register, dan apa yang terjadi kalau Mongo belum siap? - Apa yang berubah kalau
<-pollingDonedihapus? - Di mana Mike akan menambahkan publish ke Kafka, dan kenapa di sana?