06. ingest: polling

Bagian kamu: ingest/ingest.go (loop polling). Milik Grace: ingest/source.go dan klien di internal/source/bmkg dan internal/source/pvmbg. Syntax: belajar_go/07_goroutine_channel_select.go (goroutine, ticker, select), 05_interface.go, 06_context_defer.go.


Kenapa ada

Spek, bagian Model Ingest: BNPB mengambil data dengan polling berkala. Mock BMKG dan PVMBG pasif: mereka hanya menjawab saat ditanya, tidak mengirim sendiri. Intervalnya harus bisa diatur (disarankan 2 sampai 5 detik).

Kenapa dibutuhkan

Ini satu-satunya jalur data masuk ke sistem. Tanpa polling, tidak ada HazardEvent yang tersimpan, dipublikasikan, atau disajikan. Semua Problem (P1, P2, P4, P5) bergantung padanya.

Siapa mengerjakan apa

  main.go
     |  menyiapkan klien dan handler
     v
  ingest.Run   <- loop polling: ticker, cursor, goroutine per sumber   (KAMU)
     |  memanggil
     v
  SourceClient.Fetch(ctx, since)      <- antarmuka                      (GRACE)
     |  diimplementasikan oleh
     v
  klien BMKG / klien PVMBG            <- ambil data + panggil mapper    (GRACE)

Loop polling tidak tahu apa pun tentang JSON, kredensial, atau pemetaan. Ia hanya: panggil Fetch, serahkan hasilnya ke handler, geser cursor.


Bagian 1: kontrak dari Grace, SourceClient

Dikutip dari internal/ingest/source.go, baris 11 sampai 14

type SourceClient interface {
	Source() domain.Source
	Fetch(ctx context.Context, since *time.Time) ([]domain.HazardEvent, error)
}

Semua klien sumber harus punya dua method:

  • Source() menyebut namanya (BMKG atau PVMBG), untuk log.
  • Fetch(ctx, since) mengambil data yang lebih baru dari since dan mengembalikan HazardEvent yang sudah dipetakan.

since bertipe *time.Time (pointer): nil berarti “belum pernah mengambil, minta semuanya”. (Pelajaran 3 dan 5.)

Dikutip dari internal/ingest/source.go, baris 24 sampai 29

type SourceError struct {
	Source     domain.Source
	Kind       ErrorKind
	StatusCode int
	Err        error
}

Kalau Fetch gagal, ia mengembalikan SourceError yang membawa jenis kegagalan:

KindArtinyaContoh
unauthorizedkredensial ditolak401 atau 403, kunci salah
unavailablesumber tidak bisa dijangkaujaringan putus, timeout, 5xx, outage PVMBG
invalid_responsejawaban tidak bisa dipakaiJSON rusak, record gagal dipetakan

Perbedaan jenis ini dipakai loop untuk memutuskan tingkat log.


Bagian 2: Run, satu goroutine per sumber

Dikutip dari internal/ingest/ingest.go, baris 17 sampai 17

type Handler func(ctx context.Context, events []domain.HazardEvent) error

Handler adalah nama untuk sebuah bentuk fungsi: “fungsi yang menerima ctx dan daftar event, dan mengembalikan error”. main.go yang mengisinya (menyimpan ke store). Dengan cara ini ingest tidak tahu apa pun tentang store, jadi bisa diuji tanpa MongoDB. (Pelajaran 8.)

Dikutip dari internal/ingest/ingest.go, baris 21 sampai 32

func Run(ctx context.Context, logger *slog.Logger, interval time.Duration, clients []SourceClient, handle Handler) {
	var wg sync.WaitGroup
	for _, client := range clients {
		wg.Add(1)
		go func() {
			defer wg.Done()
			p := &poller{logger: logger, client: client, handle: handle}
			p.loop(ctx, interval)
		}()
	}
	wg.Wait()
}

Untuk setiap klien, dimulai satu goroutine (pekerja yang berjalan bersamaan). wg.Wait() menunggu semua goroutine selesai, yang terjadi saat service diminta berhenti.

Kenapa satu goroutine per sumber? Problem 2 menyebut head-of-line blocking: kalau BMKG dan PVMBG dipanggil berurutan dalam satu loop, BMKG yang cepat ikut menunggu PVMBG yang lambat. Dengan loop terpisah, PVMBG yang macet atau mati tidak menunda BMKG sama sekali. Terbukti: saat outage PVMBG, event BMKG tetap bertambah (bukti_p4/06_...). (Pelajaran 7.)


Bagian 3: poller dan loop

Dikutip dari internal/ingest/ingest.go, baris 34 sampai 39

type poller struct {
	logger *slog.Logger
	client SourceClient
	handle Handler
	since  *time.Time // newest OccurredAt already handled; owned by this poller's goroutine
}

poller menyimpan status satu sumber. Perhatikan since: itulah cursor, yaitu “sampai event mana saya sudah memproses”. Karena since hanya disentuh goroutine milik poller itu sendiri, tidak perlu kunci (mutex): tidak ada dua goroutine yang menulis variabel yang sama.

Dikutip dari internal/ingest/ingest.go, baris 41 sampai 53

func (p *poller) loop(ctx context.Context, interval time.Duration) {
	p.cycle(ctx)
	ticker := time.NewTicker(interval)
	defer ticker.Stop()
	for {
		select {
		case <-ctx.Done():
			return
		case <-ticker.C:
			p.cycle(ctx)
		}
	}
}

Bentuknya:

  1. p.cycle(ctx): ambil sekali langsung, tanpa menunggu satu interval dulu.
  2. Buat ticker yang berbunyi tiap interval (default 3 detik).
  3. Di dalam for, select menunggu salah satu:
    • <-ctx.Done(): service diminta berhenti (docker compose stop). Loop return.
    • <-ticker.C: waktunya polling lagi. Panggil cycle.

Kalau satu siklus lebih lama dari interval (misalnya PVMBG menggantung sampai timeout), ticker menjatuhkan tick yang terlewat, sehingga siklus tidak menumpuk.


Bagian 4: cycle, satu putaran

Dikutip dari internal/ingest/ingest.go, baris 56 sampai 83

func (p *poller) cycle(parent context.Context) {
	ctx := core.WithCorrelationID(parent, core.NewCorrelationID())
	source := string(p.client.Source())
 
	events, err := p.client.Fetch(ctx, p.since)
	if err != nil {
		if ctx.Err() != nil {
			return
		}
		p.logFetchError(ctx, source, err)
		return
	}
	if len(events) == 0 {
		p.logger.DebugContext(ctx, "poll cycle: nothing new", "source", source)
		return
	}
 
	if err := p.handle(ctx, events); err != nil {
		p.logger.ErrorContext(ctx, "handle events failed, will retry next cycle", "source", source, "error", err.Error())
		return
	}
 
	// The cursor moves only after the handler succeeded.
	if latest := newestOccurredAt(events); p.since == nil || latest.After(*p.since) {
		p.since = &latest
	}
	p.logger.InfoContext(ctx, "poll cycle done", "source", source, "events", len(events))
}

Langkah demi langkah:

  1. Correlation ID baru untuk siklus ini (WithCorrelationID(parent, NewCorrelationID())). Semua log dan panggilan keluar dalam siklus ini membawa ID yang sama.
  2. Fetch(ctx, p.since): minta data yang lebih baru dari cursor.
  3. Kalau Fetch gagal:
    • Kalau penyebabnya ctx dibatalkan (service sedang berhenti), keluar diam-diam.
    • Selain itu catat lewat logFetchError, lalu return. Cursor tidak digeser, jadi siklus berikutnya mencoba lagi dari titik yang sama.
  4. Kalau tidak ada event baru: catat di level debug dan selesai.
  5. p.handle(ctx, events): serahkan ke handler (di main.go: simpan ke Mongo). Kalau handler gagal (misalnya Mongo mati), catat ERROR dan return. Cursor tidak digeser. Ini yang disebut retry: event yang sama diambil lagi di siklus berikutnya.
  6. Hanya setelah handler berhasil, cursor maju ke OccurredAt terbaru di antara event itu, dan hanya maju, tidak pernah mundur.

Kenapa cursor hanya maju setelah berhasil? Supaya tidak ada data yang hilang diam-diam. Kalau cursor maju lebih dulu lalu penyimpanan gagal, event itu tidak akan pernah diambil lagi. Karena store.Save adalah upsert, mengambil event yang sama dua kali aman.

Dikutip dari internal/ingest/ingest.go, baris 88 sampai 99

func (p *poller) logFetchError(ctx context.Context, source string, err error) {
	attrs := []any{"source", source, "error", err.Error()}
	var se *SourceError
	if errors.As(err, &se) {
		attrs = append(attrs, "kind", string(se.Kind))
		if se.Kind == ErrorUnauthorized {
			p.logger.ErrorContext(ctx, "source rejected our credentials", attrs...)
			return
		}
	}
	p.logger.WarnContext(ctx, "poll failed", attrs...)
}

Tingkat log dibedakan menurut jenis kegagalan:

  • unauthorized → ERROR. Kredensial salah tidak akan sembuh sendiri; perlu dilihat orang.
  • selain itu → WARN. Sumber mati atau rusak itu wajar, apalagi saat demo outage PVMBG.

errors.As(err, &se) memeriksa apakah error itu berjenis *SourceError dan, kalau ya, mengisi se. (Pelajaran 2.)


Bagian 5: klien sumber (milik Grace), supaya kamu bisa menjelaskan perannya

Dikutip dari internal/source/bmkg/client.go, baris 60 sampai 92

func (client *Client) Fetch(
	ctx context.Context,
	since *time.Time,
) ([]domain.HazardEvent, error) {
	events, err := fetch[mapping.SeismicEvent](client, ctx, "/seismic-events", since)
	if err != nil {
		return nil, err
	}
	warnings, err := fetch[mapping.TsunamiWarning](client, ctx, "/tsunami-warnings", since)
	if err != nil {
		return nil, err
	}
 
	warningsByEvent := make(map[string]mapping.TsunamiWarning, len(warnings))
	for _, warning := range warnings {
		warningsByEvent[warning.RelatedEventID] = warning
	}
 
	mapped := make([]domain.HazardEvent, 0, len(events))
	for index, event := range events {
		var warning *mapping.TsunamiWarning
		if related, ok := warningsByEvent[event.EventID]; ok {
			warningCopy := related
			warning = &warningCopy
		}
		hazard, err := client.mapper.MapSeismicEvent(event, warning)
		if err != nil {
			return nil, invalidResponse(fmt.Errorf("map BMKG event at index %d: %w", index, err))
		}
		mapped = append(mapped, hazard)
	}
	return mapped, nil
}

Klien BMKG:

  1. Ambil /seismic-events, lalu /tsunami-warnings, dengan since yang sama.
  2. Indeks warning menurut related_event_id (warningsByEvent, sebuah map).
  3. Petakan tiap event bersama warning-nya kalau ada, lewat MapSeismicEvent.

Kenapa digabung? Aturan severity di spek: kalau ada TsunamiWarning, warning menimpa aturan magnitude (mis. magnitude 3 dengan warning Awas menjadi AWAS). Event dan warning datang dari dua endpoint berbeda, jadi harus digabung sebelum dipetakan.

Klien PVMBG lebih sederhana: satu endpoint, GET /volcanic-reports, lalu MapVolcanicReport untuk tiap laporan.

Kredensial berbeda dan tidak boleh tertukar: BMKG memakai header X-BMKG-Key, PVMBG memakai Authorization: Bearer <token>.

Keterbatasan klien Grace (saran, bukan bug untuk data mock; sudah diperiksa terhadap data hasil polling, 9 dari 9 gempa bertsunami memiliki warning dan severity-nya cocok):

  • Satu record yang gagal dipetakan menggagalkan seluruh Fetch, sehingga satu record rusak memblokir sumber itu terus-menerus (cursor tidak maju).
  • Warning yang terbit sesudah event-nya sudah diproses tidak mengoreksi severity event itu. Di mock, warning selalu tersedia bersamaan dengan event-nya.

Cara menjelaskan ingest dalam 30 detik

ingest.Run memberi satu goroutine per sumber, masing-masing mem-polling tiap interval. Setiap siklus membuat correlation ID baru, memanggil Fetch pada klien sumber dengan cursor since, lalu menyerahkan event ke handler yang menyimpannya. Cursor hanya maju setelah penyimpanan berhasil, jadi kegagalan di tengah jalan hanya berarti retry, bukan data hilang. Karena tiap sumber punya loop sendiri, PVMBG yang lambat atau mati tidak menunda BMKG.

Yang sudah dibuktikan

Terhadap bmkg-mock dan pvmbg-mock asli di docker compose (bukti_p4/):

  • polling berjalan dengan correlation ID per siklus;
  • outage PVMBG dicatat WARN jenis unavailable, BMKG tetap jalan, PVMBG pulih tanpa restart dan mengejar event yang terlewat;
  • perubahan skema PVMBG masuk tanpa restart;
  • Mongo mati lalu hidup: penyimpanan gagal diulang, tanpa duplikat.

Batasan

  • Cursor hanya di memori. Setelah restart, polling mengambil semua data lagi (aman karena upsert, tetapi event lama dikirim ulang ke tahap berikutnya, jadi consumer Kafka harus idempotent).
  • Mode outage hang PVMBG belum diuji, baru mode error (503).
  • ingest.go tidak punya tes unit di repo (tes yang dijalankan bersifat sementara, sesuai permintaan tidak menaruh tes di repo).

Latihan

  1. Apa yang terjadi pada cursor kalau handle mengembalikan error? Kenapa itu aman?
  2. Kenapa select memiliki case <-ctx.Done()? Apa yang terjadi saat kamu menjalankan docker compose stop bnpb-aggregator tanpa case itu?
  3. Kalau kamu menambah sumber ketiga, apa saja yang harus diubah di ingest.go? (Jawab: tidak ada. Cukup tambahkan satu klien ke daftar clients di main.go.)