04. store: Canonical Store di MongoDB

Bagian kamu (P4). File asli: internal/store/store.go. Syntax Go yang dipakai dijelaskan di belajar_go/ (nomor pelajarannya disebut di tiap bagian). Istilah yang belum familier ada di GLOSSARY.md.


Kenapa ada

Soal Problem 4 menceritakan masalah ini: HazardEvent disimpan di tabel dengan kolom tetap. Setiap kali PVMBG menambah field baru (misalnya confidence_level), tim harus mengubah tabel (migration) dan men-deploy ulang bersama-sama. Soal meminta penyimpanan yang menerima field baru tanpa migration.

Spek juga mewajibkan satu storage dimiliki satu service. Package store inilah tempat semua kode yang bicara ke MongoDB berada, dan hanya aggregator yang memakainya.

Kenapa dibutuhkan

Tanpa store, tidak ada tempat HazardEvent disimpan, dan kriteria 2 Problem 4 (“record sebelum dan sesudah confidence_level muncul tersimpan berdampingan”) tidak bisa didemokan.

Ide intinya, dalam tiga kalimat

  1. MongoDB menyimpan document (mirip JSON), bukan baris bertabel. Tidak ada schema yang harus diubah saat ada field baru.
  2. Seluruh HazardEvent disimpan apa adanya, termasuk Attributes yang isinya bebas. Jadi confidence_level yang tiba-tiba muncul langsung ikut tersimpan.
  3. Menyimpan memakai upsert (ganti kalau sudah ada, buat baru kalau belum), karena polling bisa mengambil data yang sama dua kali.

Bagian 1: struct Store

Dikutip dari internal/store/store.go, baris 18 sampai 21

type Store struct {
	client *mongo.Client
	events *mongo.Collection
}

Store hanya memegang dua hal: koneksi ke MongoDB (client) dan collection tempat event disimpan (events). Dua field itu berhuruf kecil, jadi privat: package lain tidak bisa menyentuhnya langsung, hanya lewat method milik Store. (Pelajaran 3 tentang struct dan huruf besar/kecil; pelajaran 1.)

Nama collection-nya hazard_events:

Dikutip dari internal/store/store.go, baris 16 sampai 16

const collectionName = "hazard_events"

Bagian 2: New, menyambung dan menyiapkan index

Dikutip dari internal/store/store.go, baris 25 sampai 59

func New(ctx context.Context, uri, database string) (*Store, error) {
	// Nested objects inside Attributes come back as bson.M (a map) instead of bson.D
	// (an ordered pair list), so they behave like the map[string]any they were saved from.
	client, err := mongo.Connect(options.Client().
		ApplyURI(uri).
		SetBSONOptions(&options.BSONOptions{DefaultDocumentM: true}))
	if err != nil {
		return nil, fmt.Errorf("connect mongo: %w", err)
	}
	if err := client.Ping(ctx, nil); err != nil {
		_ = client.Disconnect(ctx)
		return nil, fmt.Errorf("ping mongo: %w", err)
	}
 
	events := client.Database(database).Collection(collectionName)
	_, err = events.Indexes().CreateMany(ctx, []mongo.IndexModel{
		{
			// Natural key: the same source record is never stored twice, so a retried
			// or re-polled record is an upsert, not a duplicate.
			Keys:    bson.D{{Key: "source", Value: 1}, {Key: "source_ref_id", Value: 1}},
			Options: options.Index().SetUnique(true),
		},
		{Keys: bson.D{{Key: "ingested_at", Value: 1}}},
		// hazard_id is derived from (source, source_ref_id), so it is unique too. Serves GET by id.
		{
			Keys:    bson.D{{Key: "hazard_id", Value: 1}},
			Options: options.Index().SetUnique(true),
		},
	})
	if err != nil {
		_ = client.Disconnect(ctx)
		return nil, fmt.Errorf("create indexes: %w", err)
	}
	return &Store{client: client, events: events}, nil
}

Langkahnya berurutan:

  1. mongo.Connect(...) membuat client. Perhatikan SetBSONOptions(... DefaultDocumentM: true). Ini temuan dari pengujian: tanpa opsi itu, objek bersarang di dalam Attributes (misalnya data tsunami_warning) kembali dari Mongo bertipe bson.D, yaitu daftar pasangan berurutan, bukan map. JSON-nya masih benar, tetapi kode Go yang membacanya akan kaget. Dengan opsi ini hasilnya bson.M, yang bentuknya seperti map.
  2. client.Ping(...) memastikan Mongo benar-benar menjawab. Connect saja belum menjamin itu. Kalau Ping gagal, service gagal start dengan jelas, bukan jalan setengah-setengah lalu error di request pertama. Docker lalu menjalankannya ulang (restart: unless-stopped).
  3. CreateMany untuk indexes. Ini bukan schema. Index adalah daftar isi supaya pencarian cepat, dan tidak membatasi field apa yang boleh ada di sebuah document. Tiga index dibuat:
    • unique index pada (source, source_ref_id): satu data dari sumber (misalnya PVMBG-RPT-000030) tidak boleh tersimpan dua kali. Ini pasangan dari upsert di Save.
    • index pada ingested_at: supaya filter since tidak memeriksa seluruh collection.
    • unique index pada hazard_id: untuk GET /hazards/{id} yang cepat, dan mencegah ID ganda. hazard_id dibuat deterministik oleh NewHazardID milik Grace (hash dari source + id asli), jadi unik juga.
  4. Kalau langkah apa pun gagal, koneksi ditutup dulu (client.Disconnect) sebelum mengembalikan error, supaya tidak ada koneksi yang bocor. (Pola if err != nil di pelajaran 2; %w membungkus error.)

Bagian 3: Save, upsert

Dikutip dari internal/store/store.go, baris 62 sampai 69

func (s *Store) Save(ctx context.Context, e domain.HazardEvent) error {
	filter := bson.D{{Key: "source", Value: e.Source}, {Key: "source_ref_id", Value: e.SourceRefID}}
	_, err := s.events.ReplaceOne(ctx, filter, e, options.Replace().SetUpsert(true))
	if err != nil {
		return fmt.Errorf("save hazard event %s/%s: %w", e.Source, e.SourceRefID, err)
	}
	return nil
}
  • filter menyatakan “cari document dengan source dan source_ref_id ini”.
  • ReplaceOne(... SetUpsert(true)) berarti: kalau document yang cocok ada, ganti seluruhnya dengan e. Kalau belum ada, buat baru.
  • Kenapa ini penting: event yang sama bisa sampai dua kali ke Save (cursor hilang saat restart, penyimpanan gagal lalu di-retry, atau koreksi severity). Hasilnya selalu satu document. Operasi seperti ini disebut idempotent.
  • Di mana “tanpa migration” terjadi: e bertipe domain.HazardEvent, dan salah satu fieldnya Attributes map[string]any. Driver Mongo menulis seluruh isi struct itu apa adanya sebagai satu document. Saat PVMBG mulai mengirim confidence_level, field itu masuk ke Attributes (dikerjakan mapper Grace), dan Save menyimpannya tanpa perubahan apa pun di kode ini. (Pelajaran 4 tentang map[string]any.)

Bagian 4: Filter dan List, membaca dengan syarat

Dikutip dari internal/store/store.go, baris 72 sampai 77

type Filter struct {
	Since      *time.Time // events ingested strictly after this instant
	Source     domain.Source
	HazardType domain.HazardType
	Severity   domain.Severity
}

Filter adalah kumpulan syarat. Nilai kosong berarti “tanpa syarat ini”. Since *time.Time berupa pointer supaya bisa dibedakan antara “tidak ada batas waktu” (nil) dan “ada batas waktu” (isi). time.Time biasa selalu punya nilai, jadi tidak bisa “kosong”. (Pelajaran 3, bagian pointer.)

Dikutip dari internal/store/store.go, baris 80 sampai 103

func (s *Store) List(ctx context.Context, f Filter) ([]domain.HazardEvent, error) {
	filter := bson.D{}
	if f.Since != nil {
		filter = append(filter, bson.E{Key: "ingested_at", Value: bson.D{{Key: "$gt", Value: *f.Since}}})
	}
	if f.Source != "" {
		filter = append(filter, bson.E{Key: "source", Value: f.Source})
	}
	if f.HazardType != "" {
		filter = append(filter, bson.E{Key: "hazard_type", Value: f.HazardType})
	}
	if f.Severity != "" {
		filter = append(filter, bson.E{Key: "severity", Value: f.Severity})
	}
	cur, err := s.events.Find(ctx, filter, options.Find().SetSort(bson.D{{Key: "ingested_at", Value: 1}}))
	if err != nil {
		return nil, fmt.Errorf("find hazard events: %w", err)
	}
	events := []domain.HazardEvent{}
	if err := cur.All(ctx, &events); err != nil {
		return nil, fmt.Errorf("decode hazard events: %w", err)
	}
	return events, nil
}

Cara kerjanya:

  1. Mulai dari filter kosong bson.D{}. Itu artinya “cocok dengan semua”.
  2. Untuk setiap syarat yang diisi pemanggil, tambahkan satu pasangan ke filter. Mongo menggabungkan semuanya dengan AND (semua harus terpenuhi). $gt berarti “lebih besar dari”, jadi since bersifat eksklusif: event dengan ingested_at persis sama tidak ikut.
  3. Find(...SetSort(ingested_at: 1)) menjalankan query dan mengurutkan dari yang paling awal masuk.
  4. cur.All(ctx, &events) membaca semua hasil dan mengubah tiap document menjadi HazardEvent. Document lama (tanpa confidence_level) dan baru (dengan itu) sama-sama terbaca. Inilah bukti kriteria 2.
  5. events := []domain.HazardEvent{} sengaja berupa slice kosong, bukan nil. Akibatnya JSON yang dihasilkan [], bukan null, sehingga API menjawab {"data":[]} untuk hasil kosong. (Pelajaran 4.)

Kenapa ingested_at, bukan occurred_at: ingested_at adalah kapan BNPB menerima data. Pembaca yang bertanya “apa yang baru sejak terakhir saya cek?” tidak akan melewatkan event yang datang terlambat, walau kejadiannya sudah lama.


Bagian 5: Get, satu event menurut hazard_id

Dikutip dari internal/store/store.go, baris 106 sampai 116

func (s *Store) Get(ctx context.Context, hazardID string) (*domain.HazardEvent, error) {
	var e domain.HazardEvent
	err := s.events.FindOne(ctx, bson.D{{Key: "hazard_id", Value: hazardID}}).Decode(&e)
	if errors.Is(err, mongo.ErrNoDocuments) {
		return nil, nil
	}
	if err != nil {
		return nil, fmt.Errorf("get hazard event %s: %w", hazardID, err)
	}
	return &e, nil
}
  • Mengembalikan pointer *domain.HazardEvent. nil berarti “tidak ada”.
  • “Tidak ada” bukan kegagalan. Karena itu mongo.ErrNoDocuments diubah menjadi return nil, nil (hasil kosong, error kosong). Pemanggil (api) lalu menjawab 404. Kalau Mongo benar-benar bermasalah, errornya dikembalikan dan api menjawab 500. (Pelajaran 2: errors.Is.)

Cara menjelaskan store dalam 30 detik

store adalah satu-satunya tempat yang bicara ke MongoDB. Setiap HazardEvent disimpan sebagai satu document tanpa schema, jadi field baru dari sumber, seperti confidence_level, langsung tersimpan tanpa migration. Menyimpannya memakai upsert dengan kunci (source, source_ref_id), sehingga data yang sama yang terambil dua kali tidak menjadi dobel. Pembacaan lewat List dengan filter dan Get menurut hazard_id, dan keduanya hanya dipanggil oleh API internal aggregator.

Yang sudah dibuktikan

Di docker compose dengan Mongo asli (bukti_p4/04_kriteria2_atribut_dinamis.txt): 31 record PVMBG tanpa confidence_level tetap terbaca setelah skema berubah, dan 2 record baru dengan confidence_level ada di collection yang sama. Mongo dimatikan lalu dinyalakan: tidak ada duplikat (bukti_p4/07_...).

Batasan

  • Satu node MongoDB tanpa replikasi.
  • Aggregator memakai akun root MongoDB (hak akses minimum belum diterapkan).
  • List tanpa pagination: mengembalikan semua yang cocok.
  • Nilai bersarang di Attributes bertipe bson.M, sehingga .(map[string]any) gagal padanya. Kode yang membaca isi bersarang harus memakai bson.M atau lewat JSON.

Latihan

  1. Kalau Save memakai InsertOne dan bukan ReplaceOne(... upsert), apa yang terjadi saat polling mengambil event yang sama dua kali? (Index unik akan menolak yang kedua, lalu handle gagal terus.)
  2. Kenapa List mengembalikan []domain.HazardEvent{} dan bukan nil saat tidak ada hasil?
  3. Kenapa Get mengembalikan (nil, nil) untuk “tidak ada” dan bukan sebuah error?