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
- MongoDB menyimpan document (mirip JSON), bukan baris bertabel. Tidak ada schema yang harus diubah saat ada field baru.
- Seluruh
HazardEventdisimpan apa adanya, termasukAttributesyang isinya bebas. Jadiconfidence_levelyang tiba-tiba muncul langsung ikut tersimpan. - 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:
mongo.Connect(...)membuat client. PerhatikanSetBSONOptions(... DefaultDocumentM: true). Ini temuan dari pengujian: tanpa opsi itu, objek bersarang di dalamAttributes(misalnya datatsunami_warning) kembali dari Mongo bertipebson.D, yaitu daftar pasangan berurutan, bukanmap. JSON-nya masih benar, tetapi kode Go yang membacanya akan kaget. Dengan opsi ini hasilnyabson.M, yang bentuknya sepertimap.client.Ping(...)memastikan Mongo benar-benar menjawab.Connectsaja 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).CreateManyuntuk 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 (misalnyaPVMBG-RPT-000030) tidak boleh tersimpan dua kali. Ini pasangan dari upsert diSave. - index pada
ingested_at: supaya filtersincetidak memeriksa seluruh collection. - unique index pada
hazard_id: untukGET /hazards/{id}yang cepat, dan mencegah ID ganda.hazard_iddibuat deterministik olehNewHazardIDmilik Grace (hash dari source + id asli), jadi unik juga.
- unique index pada
- Kalau langkah apa pun gagal, koneksi ditutup dulu (
client.Disconnect) sebelum mengembalikan error, supaya tidak ada koneksi yang bocor. (Polaif err != nildi pelajaran 2;%wmembungkus 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
}filtermenyatakan “cari document dengansourcedansource_ref_idini”.ReplaceOne(... SetUpsert(true))berarti: kalau document yang cocok ada, ganti seluruhnya dengane. 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:
ebertipedomain.HazardEvent, dan salah satu fieldnyaAttributes map[string]any. Driver Mongo menulis seluruh isi struct itu apa adanya sebagai satu document. Saat PVMBG mulai mengirimconfidence_level, field itu masuk keAttributes(dikerjakan mapper Grace), danSavemenyimpannya tanpa perubahan apa pun di kode ini. (Pelajaran 4 tentangmap[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:
- Mulai dari filter kosong
bson.D{}. Itu artinya “cocok dengan semua”. - Untuk setiap syarat yang diisi pemanggil, tambahkan satu pasangan ke filter. Mongo
menggabungkan semuanya dengan AND (semua harus terpenuhi).
$gtberarti “lebih besar dari”, jadisincebersifat eksklusif: event denganingested_atpersis sama tidak ikut. Find(...SetSort(ingested_at: 1))menjalankan query dan mengurutkan dari yang paling awal masuk.cur.All(ctx, &events)membaca semua hasil dan mengubah tiap document menjadiHazardEvent. Document lama (tanpaconfidence_level) dan baru (dengan itu) sama-sama terbaca. Inilah bukti kriteria 2.events := []domain.HazardEvent{}sengaja berupa slice kosong, bukannil. Akibatnya JSON yang dihasilkan[], bukannull, 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.nilberarti “tidak ada”. - “Tidak ada” bukan kegagalan. Karena itu
mongo.ErrNoDocumentsdiubah menjadireturn nil, nil(hasil kosong, error kosong). Pemanggil (api) lalu menjawab 404. Kalau Mongo benar-benar bermasalah, errornya dikembalikan danapimenjawab 500. (Pelajaran 2:errors.Is.)
Cara menjelaskan store dalam 30 detik
storeadalah satu-satunya tempat yang bicara ke MongoDB. Setiap HazardEvent disimpan sebagai satu document tanpa schema, jadi field baru dari sumber, seperticonfidence_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 lewatListdengan filter danGetmenuruthazard_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).
Listtanpa pagination: mengembalikan semua yang cocok.- Nilai bersarang di
Attributesbertipebson.M, sehingga.(map[string]any)gagal padanya. Kode yang membaca isi bersarang harus memakaibson.Matau lewat JSON.
Latihan
- Kalau
SavememakaiInsertOnedan bukanReplaceOne(... upsert), apa yang terjadi saat polling mengambil event yang sama dua kali? (Index unik akan menolak yang kedua, laluhandlegagal terus.) - Kenapa
Listmengembalikan[]domain.HazardEvent{}dan bukannilsaat tidak ada hasil? - Kenapa
Getmengembalikan(nil, nil)untuk “tidak ada” dan bukan sebuah error?