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)
	}
}
  1. core.RunHealthcheckIfRequested(): kalau program dijalankan sebagai /app healthcheck (oleh HEALTHCHECK Docker), ia memeriksa dirinya sendiri lalu keluar. Kalau tidak, tidak berefek. (Lihat 03_pondasi_core.md.)
  2. config.Load(): baca konfigurasi. Kalau ada variabel wajib yang kosong, semua masalah dicetak sekaligus lalu keluar dengan kode 1.
  3. 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(...) dan pvmbg.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:

  1. stop() memastikan ctx batal, juga kalau Serve berhenti karena error.
  2. <-pollingDone menunggu poller benar-benar berhenti.
  3. return serveErr, lalu defer dijalankan: 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

  1. Docker mengirim SIGTERM.
  2. ctx dibatalkan.
  3. Serve berhenti menerima koneksi baru, memberi request yang berjalan sampai 10 detik.
  4. ingest.Run berhenti: panggilan HTTP yang berjalan dibatalkan, loop memantau ctx.
  5. <-pollingDone terbuka.
  6. defer menutup koneksi Mongo.
  7. 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.go merakit 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

  1. Kenapa store.New dipanggil sebelum api.Register, dan apa yang terjadi kalau Mongo belum siap?
  2. Apa yang berubah kalau <-pollingDone dihapus?
  3. Di mana Mike akan menambahkan publish ke Kafka, dan kenapa di sana?