10 months ago
I'm trying to set up a multi-region replica with MongoDB Atlas using multi-region electable nodes.
I'm using Golang with MongoDB.
I am trying to use Railway with 3 replicas:
- California
- Netherlands
- Singapore
With MongoDB Atlas on GCP in the same locations as above, I expected that each region would read from its nearest database instance.
However, the issue is that on the Singapore Railway server, most of the time it reads from the MongoDB California server.
This makes my test API, which has x10 read queries, much slower — from 300ms to 3s.
My question is: Do I need to do any additional setup to make this work, or are there some service limitations I should be aware of?
func mapRegionToTags(railwayRegion string) []tag.Set {
if strings.HasPrefix(railwayRegion, "asia-southeast1") {
return []tag.Set{
{{Name: "region", Value: "SOUTHEASTERN_ASIA_PACIFIC"}}, // 1st
{{Name: "region", Value: "EUROPE_WEST_4"}}, // 2nd
{}, // 3rd (nearest)
}
}
if strings.HasPrefix(railwayRegion, "europe-west4") {
return []tag.Set{
{{Name: "region", Value: "EUROPE_WEST_4"}},
{{Name: "region", Value: "SOUTHEASTERN_ASIA_PACIFIC"}},
{},
}
}
if strings.HasPrefix(railwayRegion, "us-west") {
return []tag.Set{
{{Name: "region", Value: "US_WEST_2"}},
{{Name: "region", Value: "EUROPE_WEST_4"}},
{},
}
}
// Fallback
return []tag.Set{{}}
}
func SetUp() {
//TODO: may have to find tune the keep alive, timeout duration to make the best one
dialer := &net.Dialer{
KeepAlive: 300 * time.Second,
}
rawRegion := os.Getenv("RAILWAY_REPLICA_REGION")
tagSets := mapRegionToTags(rawRegion)
fmt.Println("Region map: ", rawRegion, tagSets)
rp, err := readpref.New(
readpref.NearestMode,
readpref.WithTagSets(tagSets...),
)
if err != nil {
// If we can't create the preference, we shouldn't start the server.
log.Fatalf("Failed to create MongoDB ReadPreference: %v", err)
}
wc := writeconcern.New(writeconcern.W(1))
clientOptions := options.Client().ApplyURI(os.Getenv("DATABASE_URL")).SetMinPoolSize(5).SetDialer(dialer).SetWriteConcern(wc).SetReadPreference(rp)
client, err := mongo.Connect(context.TODO(), clientOptions)
if err != nil {
log.Fatal(err)
}
Client = client
}15 Replies
10 months ago
Hey there! We've found the following might help you get unblocked faster:
If you find the answer from one of these, please let us know by solving the thread!
10 months ago
This thread has been marked as public for community involvement, as it does not contain any sensitive or personal information. Any further activity in this thread will be visible to everyone.
Status changed to Open brody • 11 months ago
10 months ago
If I understand your setup correctly, you have 3 MongoDB Atlas regions and you have 3 Railway regions and you want each railway replica to the appropriate mongo replica but the issue is that the singapore railway replica most often uses the california mongo replica, is my understanding correct? If so, does that behaviour happen with your current code?
My current thinking is you may have maxStalenessSeconds configured somewhere which might result in the first 2 regions being avoided if its lagging behind the primary which is especially possible if the cali region holds your primary replica
10 months ago
what it seems like is the singapore process is hitting the fallback branch, and the nearest ends up selecting California. you can wind up debugging it by logigng the rawRegion and tagSets, and hardcode a singapore preference to confirm routing.
could you also check your DATABASE_URL? a read preference in the URL will override the Go option and can get it to read from the primary in California
dev
If I understand your setup correctly, you have 3 MongoDB Atlas regions and you have 3 Railway regions and you want each railway replica to the appropriate mongo replica but the issue is that the singapore railway replica most often uses the california mongo replica, is my understanding correct? If so, does that behaviour happen with your current code? My current thinking is you may have `maxStalenessSeconds` configured somewhere which might result in the first 2 regions being avoided if its lagging behind the primary which is especially possible if the cali region holds your primary replica
10 months ago
FYI:
in mongoDB atlas I used
M10 with Multi-Cloud, Multi-Region
The primary node is California > Netherlands > Singapore
Yes, from my inspection singapore railway replica most often uses the california mongo replica.
But I'm also newbie in the horizontal scalling as well,
yeeet
what it seems like is the singapore process is hitting the fallback branch, and the nearest ends up selecting California. you can wind up debugging it by logigng the rawRegion and tagSets, and hardcode a singapore preference to confirm routing. could you also check your `DATABASE_URL`? a read preference in the URL will override the Go option and can get it to read from the primary in California
10 months ago
in my DATABASE_URL it have
/?retryWrites=true&w=majority
here
But how to hardcode a singapore preference to confirm routing ?
in my DATABASE\_URL it have /?retryWrites=true&w=majority here But how to hardcode a singapore preference to confirm routing ?
10 months ago
I'd use this for testing purposes, but it would look something like this:
rawRegion := os.Getenv("RAILWAY_REPLICA_REGION")
log.Printf("Hardcoding Singapore read preference, RAILWAY_REPLICA_REGION=%s", rawRegion)
rp, err := readpref.New(
readpref.NearestMode,
readpref.WithTagSets(
tag.Set{{Name: "region", Value: "SOUTHEASTERN_ASIA_PACIFIC"}},
),
)
if err != nil {
log.Fatalf("Singapore read preference error: %v", err)
}
clientOptions := options.Client().
ApplyURI(os.Getenv("DATABASE_URL")).
SetReadPreference(rp)9 months ago
the behavior you’re seeing is expected with your current setup. nearest does not mean “same region”, it means “lowest ping among eligible nodes”. since your primary is in california, it’s often within the latency window and gets picked, especially if the singapore secondary is slightly lagging
key points to fix / verify:
- atlas already auto-tags nodes with
region, you don’t need custom tags - your code is fine only if
RAILWAY_REPLICA_REGIONreally matchesasia-southeast1, otherwise you silently hit the{}fallback and nearest = california - log
rawRegion+tagSetson startup to confirm you’re not falling back - temporarily hardcode
{region:SOUTHEASTERN_ASIA_PACIFIC}to confirm routing (good test) - make sure you are not setting
maxStalenessSecondsanywhere — that will disqualify the singapore secondary and force reads to california - if you want to avoid the primary completely, use
secondaryPreferredinstead ofnearest
dev
If I understand your setup correctly, you have 3 MongoDB Atlas regions and you have 3 Railway regions and you want each railway replica to the appropriate mongo replica but the issue is that the singapore railway replica most often uses the california mongo replica, is my understanding correct? If so, does that behaviour happen with your current code? My current thinking is you may have `maxStalenessSeconds` configured somewhere which might result in the first 2 regions being avoided if its lagging behind the primary which is especially possible if the cali region holds your primary replica
9 months ago
no railway limitation here , this is mongo read preference behavior + latency + possible staleness filtering
a month ago
The read-preference diagnosis in this thread is right, and the tagSets fix
will work — but there's a Railway-specific wall behind it that'll bite you
the moment you try to implement it: replicas of a single Railway service all
share one environment, which means one MONGO_URI. You can't give your
Singapore replica different readPreferenceTags than your California replica
if they're replicas of the same service — there's nowhere to put a
per-replica connection string.
What actually works on Railway: deploy three separate services, one per
region, each with its own MONGO_URI carrying its own
readPreferenceTags=region:AP_SOUTHEAST_1 (or whatever Atlas's tag values are
in your cluster — check via the Atlas UI or rs.conf()), sitting behind
whatever routing you're already using to reach the right region. More
services to manage, but it's the only way to get deterministic region-local
reads with per-instance config on this platform.
Two things to set expectations on regardless of which read-preference mode
you land on: writes always go to the primary no matter what, so if your 3s
number includes writes that part won't move; and if you go with plain
nearest/secondaryPreferred instead of tags, that's accepting eventual
consistency on those reads, worth confirming is fine for your x10 query set
before shipping it.
6 days ago
One Railway service can do this: each replica gets its own RAILWAY_REPLICA_REGION, so your Go process can select its read preference at startup while all replicas share DATABASE_URL. Three separate services are not required just to vary read routing. Railway's variable reference
Your Singapore, Netherlands and Los Angeles tag names are valid GCP Atlas region names. Keep SOUTHEASTERN_ASIA_PACIFIC for this GCP cluster; AP_SOUTHEAST_1 is an AWS example. Match the tags actually advertised by your nodes. Atlas GCP regions
The important diagnostic is your ordered fallback. With a healthy eligible Singapore node matching the first tag set, nearest filters to that set before comparing latency. California being primary does not override the Singapore tag. With your configuration, California becomes eligible via the final {} if neither earlier set matches eligible nodes. Check the actual node tags/availability and which read preference the operation uses. Atlas tag selection
For a temporary Singapore-only diagnostic, keep your URI and existing client setup, but replace the read preference with:
rp, err := readpref.New(
readpref.NearestMode,
readpref.WithTags(
"provider", "GCP",
"region", "SOUTHEASTERN_ASIA_PACIFIC",
),
)
if err != nil {
return err // adapt to your setup function's error handling
}
clientOptions := options.Client().
ApplyURI(os.Getenv("DATABASE_URL")).
SetReadPreference(rp).
SetServerSelectionTimeout(10 * time.Second)This intentionally removes both fallback sets. If no eligible matching node is discoverable, a read times out instead of silently routing elsewhere. Use it to diagnose a representative read before choosing your production fallback policy.
Attach a ServerMonitor.ServerDescriptionChanged callback before connecting and log NewDescription.Addr, Kind, Tags, and AverageRTT. A CommandMonitor.Started callback can log CommandName and ConnectionID for find/aggregate to identify which host actually receives the read; logging credentials or query bodies is unnecessary. Also log the raw Railway region and final read preference, and check whether the database/collection/session overrides the client preference.
Two corrections to earlier suggestions: in your shown call order, SetReadPreference comes after ApplyURI, so that setter takes precedence. And secondaryPreferred still permits primary fallback; it does not categorically exclude California. Writes still target the primary of this replica set.
I verified these option and selection behaviors with Go driver v1.17.10 tests using synthetic replica-set descriptions, including ordered fallback, strict-local failure, and reversed option order as a control. This does not verify your live topology or identify its root cause. If the strict read still reaches California, inspect the effective preference on that operation rather than splitting the Railway service.
5 days ago
package mongo
import (
"context"
"fmt"
"log"
"net"
"os"
"time"
"go.mongodb.org/mongo-driver/mongo"
"go.mongodb.org/mongo-driver/mongo/options"
"go.mongodb.org/mongo-driver/mongo/readpref"
"go.mongodb.org/mongo-driver/mongo/writeconcern"
"replica/internal/diag")
var (
Client *mongo.Client
regionDiag *diag.RegionDiag
diagCancel context.CancelFunc)
func SetUp() {
//TODO: may have to find tune the keep alive, timeout duration to make the best one
dialer := &net.Dialer{
KeepAlive: 300 * time.Second,
}
rawRegion := os.Getenv("RAILWAY_REPLICA_REGION")
tagSets := diag.MapRegionToTags(rawRegion)
fmt.Println("Region map: ", rawRegion, tagSets)
rp, err := readpref.New(
readpref.NearestMode,
readpref.WithTagSets(tagSets...),
)
if err != nil {
// If we can't create the preference, we shouldn't start the server.
log.Fatalf("Failed to create MongoDB ReadPreference: %v", err)
}
regionDiag = diag.NewRegionDiag()
wc := writeconcern.New(writeconcern.W(1))
clientOptions := options.Client().ApplyURI(os.Getenv("DATABASE_URL")).SetMinPoolSize(5).SetDialer(dialer).SetWriteConcern(wc).SetReadPreference(rp)
clientOptions.SetServerMonitor(regionDiag.ServerMonitor())
clientOptions.SetMonitor(regionDiag.CommandMonitor())
client, err := mongo.Connect(context.TODO(), clientOptions)
if err != nil {
log.Fatal(err)
}
Client = client
diagCtx, cancelDiag := context.WithCancel(context.Background())
diagCancel = cancelDiag
go regionDiag.Report(diagCtx, client, rp)
regionDiag.StartRouteLogger(diagCtx, 60*time.Second)}
get-them666
package mongo import ( "context" "fmt" "log" "net" "os" "time" "go.mongodb.org/mongo-driver/mongo" "go.mongodb.org/mongo-driver/mongo/options" "go.mongodb.org/mongo-driver/mongo/readpref" "go.mongodb.org/mongo-driver/mongo/writeconcern" "replica/internal/diag" ) var ( Client *mongo.Client regionDiag *diag.RegionDiag diagCancel context.CancelFunc ) func SetUp() { //TODO: may have to find tune the keep alive, timeout duration to make the best one dialer := &net.Dialer{ KeepAlive: 300 * time.Second, } rawRegion := os.Getenv("RAILWAY_REPLICA_REGION") tagSets := diag.MapRegionToTags(rawRegion) fmt.Println("Region map: ", rawRegion, tagSets) rp, err := readpref.New( readpref.NearestMode, readpref.WithTagSets(tagSets...), ) if err != nil { // If we can't create the preference, we shouldn't start the server. log.Fatalf("Failed to create MongoDB ReadPreference: %v", err) } regionDiag = diag.NewRegionDiag() wc := writeconcern.New(writeconcern.W(1)) clientOptions := options.Client().ApplyURI(os.Getenv("DATABASE_URL")).SetMinPoolSize(5).SetDialer(dialer).SetWriteConcern(wc).SetReadPreference(rp) clientOptions.SetServerMonitor(regionDiag.ServerMonitor()) clientOptions.SetMonitor(regionDiag.CommandMonitor()) client, err := mongo.Connect(context.TODO(), clientOptions) if err != nil { log.Fatal(err) } Client = client diagCtx, cancelDiag := context.WithCancel(context.Background()) diagCancel = cancelDiag go regionDiag.Report(diagCtx, client, rp) regionDiag.StartRouteLogger(diagCtx, 60*time.Second) }
5 days ago
package diag
import (
"context"
"fmt"
"log"
"os"
"sort"
"strings"
"sync"
"time"
"go.mongodb.org/mongo-driver/bson"
"go.mongodb.org/mongo-driver/event"
"go.mongodb.org/mongo-driver/mongo"
"go.mongodb.org/mongo-driver/mongo/description"
"go.mongodb.org/mongo-driver/mongo/readpref"
"go.mongodb.org/mongo-driver/tag")
type RegionDiag struct {
mu sync.Mutex
topology description.Topology
haveTopo bool
routeMu sync.Mutex
routeHosts map[string]int
routeOps int}
func NewRegionDiag() *RegionDiag {
return &RegionDiag{routeHosts: make(map[string]int)}}
func (d *RegionDiag) ServerMonitor() *event.ServerMonitor {
return &event.ServerMonitor{
TopologyDescriptionChanged: func(e *event.TopologyDescriptionChangedEvent) {
d.mu.Lock()
d.topology = e.NewDescription
d.haveTopo = true
d.mu.Unlock()
},
}}
func (d *RegionDiag) CommandMonitor() *event.CommandMonitor {
return &event.CommandMonitor{
Started: func(_ context.Context, e *event.CommandStartedEvent) {
d.routeMu.Lock()
d.routeOps++
d.routeHosts[e.ConnectionID]++
d.routeMu.Unlock()
},
}}
func (d *RegionDiag) snapshot() (description.Topology, bool) {
d.mu.Lock()
defer d.mu.Unlock()
return d.topology, d.haveTopo}
func (d *RegionDiag) StartRouteLogger(ctx context.Context, every time.Duration) {
go func() {
t := time.NewTicker(every)
defer t.Stop()
for {
select {
case <-ctx.Done():
return
case <-t.C:
d.logRouteSummary()
}
}
}()}
func (d *RegionDiag) SampleReads(ctx context.Context, client *mongo.Client, n int) {
admin := client.Database("admin")
var out bson.M
err := admin.RunCommand(ctx, bson.D{{Key: "listDatabases", Value: 1}}).Decode(&out)
if err != nil {
log.Println("sample reads: could not list databases (", err, ")")
return
}
dbs, _ := out["databases"].(bson.A)
if len(dbs) == 0 {
log.Println("sample reads: no databases visible to this user")
return
}
dbName := "admin"
for _, d0 := range dbs {
if m, ok := d0.(bson.D); ok {
for _, e := range m {
if e.Key == "name" {
if s, ok := e.Value.(string); ok && s != "admin" && s != "config" && s != "local" {
dbName = s
}
}
}
}
}
coll, err := admin.Collection("__diag_probe").EstimatedDocumentCount(ctx)
if err != nil {
coll = -1
}
log.Printf("sample reads: running %d reads against %s (probe coll estimate=%d)", n, dbName, coll)
start := time.Now()
for i := 0; i < n; i++ {
err := admin.RunCommand(ctx, bson.D{{Key: "ping", Value: 1}}).Err()
if err != nil {
log.Println(" read", i, "failed:", err)
continue
}
}
log.Printf("sample reads: %d ops in %s", n, time.Since(start).Round(time.Millisecond))}
func MapRegionToTags(railwayRegion string) []tag.Set {
if strings.HasPrefix(railwayRegion, "asia-southeast1") {
return []tag.Set{
{{Name: "region", Value: "SOUTHEASTERN_ASIA_PACIFIC"}},
{{Name: "region", Value: "EUROPE_WEST_4"}},
{},
}
}
if strings.HasPrefix(railwayRegion, "europe-west4") {
return []tag.Set{
{{Name: "region", Value: "EUROPE_WEST_4"}},
{{Name: "region", Value: "SOUTHEASTERN_ASIA_PACIFIC"}},
{},
}
}
if strings.HasPrefix(railwayRegion, "us-west") {
return []tag.Set{
{{Name: "region", Value: "US_WEST_2"}},
{{Name: "region", Value: "EUROPE_WEST_4"}},
{},
}
}
return []tag.Set{{}}}
func (d *RegionDiag) LogRouteSummary() {
d.logRouteSummary()}
func (d *RegionDiag) logRouteSummary() {
d.routeMu.Lock()
hosts := make(map[string]int, len(d.routeHosts))
for k, v := range d.routeHosts {
hosts[k] = v
}
total := d.routeOps
d.routeMu.Unlock()
if total == 0 {
return
}
topo, _ := d.snapshot()
regionByHost := make(map[string]string)
for _, s := range topo.Servers {
regionByHost[s.Addr.String()] = tagValue(s.Tags, "region")
}
type row struct {
host string
region string
n int
}
rows := make([]row, 0, len(hosts))
for h, n := range hosts {
r := regionByHost[h]
if r == "" {
r = "unknown"
}
rows = append(rows, row{h, r, n})
}
sort.Slice(rows, func(i, j int) bool { return rows[i].n > rows[j].n })
var b strings.Builder
fmt.Fprintf(&b, "route summary: %d ops | replica=%s", total, replicaRegion())
for _, r := range rows {
pct := 100 * r.n / total
fmt.Fprintf(&b, " | %s [%s] %d (%d%%)", r.host, r.region, r.n, pct)
}
log.Println(b.String())}
get-them666
package diag import ( "context" "fmt" "log" "os" "sort" "strings" "sync" "time" "go.mongodb.org/mongo-driver/bson" "go.mongodb.org/mongo-driver/event" "go.mongodb.org/mongo-driver/mongo" "go.mongodb.org/mongo-driver/mongo/description" "go.mongodb.org/mongo-driver/mongo/readpref" "go.mongodb.org/mongo-driver/tag" ) type RegionDiag struct { mu sync.Mutex topology description.Topology haveTopo bool routeMu sync.Mutex routeHosts map[string]int routeOps int } func NewRegionDiag() *RegionDiag { return &RegionDiag{routeHosts: make(map[string]int)} } func (d *RegionDiag) ServerMonitor() *event.ServerMonitor { return &event.ServerMonitor{ TopologyDescriptionChanged: func(e *event.TopologyDescriptionChangedEvent) { d.mu.Lock() d.topology = e.NewDescription d.haveTopo = true d.mu.Unlock() }, } } func (d *RegionDiag) CommandMonitor() *event.CommandMonitor { return &event.CommandMonitor{ Started: func(_ context.Context, e *event.CommandStartedEvent) { d.routeMu.Lock() d.routeOps++ d.routeHosts[e.ConnectionID]++ d.routeMu.Unlock() }, } } func (d *RegionDiag) snapshot() (description.Topology, bool) { d.mu.Lock() defer d.mu.Unlock() return d.topology, d.haveTopo } func (d *RegionDiag) StartRouteLogger(ctx context.Context, every time.Duration) { go func() { t := time.NewTicker(every) defer t.Stop() for { select { case <-ctx.Done(): return case <-t.C: d.logRouteSummary() } } }() } func (d *RegionDiag) SampleReads(ctx context.Context, client *mongo.Client, n int) { admin := client.Database("admin") var out bson.M err := admin.RunCommand(ctx, bson.D{{Key: "listDatabases", Value: 1}}).Decode(&out) if err != nil { log.Println("sample reads: could not list databases (", err, ")") return } dbs, _ := out["databases"].(bson.A) if len(dbs) == 0 { log.Println("sample reads: no databases visible to this user") return } dbName := "admin" for _, d0 := range dbs { if m, ok := d0.(bson.D); ok { for _, e := range m { if e.Key == "name" { if s, ok := e.Value.(string); ok && s != "admin" && s != "config" && s != "local" { dbName = s } } } } } coll, err := admin.Collection("__diag_probe").EstimatedDocumentCount(ctx) if err != nil { coll = -1 } log.Printf("sample reads: running %d reads against %s (probe coll estimate=%d)", n, dbName, coll) start := time.Now() for i := 0; i < n; i++ { err := admin.RunCommand(ctx, bson.D{{Key: "ping", Value: 1}}).Err() if err != nil { log.Println(" read", i, "failed:", err) continue } } log.Printf("sample reads: %d ops in %s", n, time.Since(start).Round(time.Millisecond)) } func MapRegionToTags(railwayRegion string) []tag.Set { if strings.HasPrefix(railwayRegion, "asia-southeast1") { return []tag.Set{ {{Name: "region", Value: "SOUTHEASTERN_ASIA_PACIFIC"}}, {{Name: "region", Value: "EUROPE_WEST_4"}}, {}, } } if strings.HasPrefix(railwayRegion, "europe-west4") { return []tag.Set{ {{Name: "region", Value: "EUROPE_WEST_4"}}, {{Name: "region", Value: "SOUTHEASTERN_ASIA_PACIFIC"}}, {}, } } if strings.HasPrefix(railwayRegion, "us-west") { return []tag.Set{ {{Name: "region", Value: "US_WEST_2"}}, {{Name: "region", Value: "EUROPE_WEST_4"}}, {}, } } return []tag.Set{{}} } func (d *RegionDiag) LogRouteSummary() { d.logRouteSummary() } func (d *RegionDiag) logRouteSummary() { d.routeMu.Lock() hosts := make(map[string]int, len(d.routeHosts)) for k, v := range d.routeHosts { hosts[k] = v } total := d.routeOps d.routeMu.Unlock() if total == 0 { return } topo, _ := d.snapshot() regionByHost := make(map[string]string) for _, s := range topo.Servers { regionByHost[s.Addr.String()] = tagValue(s.Tags, "region") } type row struct { host string region string n int } rows := make([]row, 0, len(hosts)) for h, n := range hosts { r := regionByHost[h] if r == "" { r = "unknown" } rows = append(rows, row{h, r, n}) } sort.Slice(rows, func(i, j int) bool { return rows[i].n > rows[j].n }) var b strings.Builder fmt.Fprintf(&b, "route summary: %d ops | replica=%s", total, replicaRegion()) for _, r := range rows { pct := 100 * r.n / total fmt.Fprintf(&b, " | %s [%s] %d (%d%%)", r.host, r.region, r.n, pct) } log.Println(b.String()) }
5 days ago
continued
func (d *RegionDiag) Report(ctx context.Context, client *mongo.Client, rp *readpref.ReadPref) {
log.Println("========== MONGO REGION DIAGNOSTIC ==========")
log.Println("replica id: ", os.Getenv("RAILWAY_REPLICA_ID"))
log.Println("railway region: ", replicaRegion())
log.Println("DATABASE_URL host:", uriHost(os.Getenv("DATABASE_URL")))
if rp != nil {
log.Println("read pref mode: ", rp.Mode())
if ms, ok := rp.MaxStaleness(); ok {
log.Println("maxStaleness: ", ms, " <-- SET, this can disqualify remote secondaries")
} else {
log.Println("maxStaleness: <unset>")
}
tagSets := rp.TagSets()
for i, ts := range tagSets {
log.Printf("tag set[%d]: %s", i, formatTagSet(ts))
}
}
var pingErr error
if err := client.Ping(ctx, nil); err != nil {
pingErr = err
log.Println("ping: FAILED:", err)
} else {
log.Println("ping: ok")
}
time.Sleep(1500 * time.Millisecond)
topo, ok := d.snapshot()
if !ok {
log.Println("topology: no topology description captured yet")
} else {
log.Println("topology kind: ", topo.Kind.String())
d.reportServers(topo)
d.reportRegions(topo)
d.reportTagEvaluation(topo, rp)
}
if pingErr == nil {
d.reportReplicaSetStatus(ctx, client)
}
log.Println("========== END DIAGNOSTIC ==========")}
func (d *RegionDiag) reportServers(topo description.Topology) {
if len(topo.Servers) == 0 {
log.Println("servers: none discovered")
return
}
log.Println("servers discovered by driver:", len(topo.Servers))
servers := append([]description.Server(nil), topo.Servers...)
sort.Slice(servers, func(i, j int) bool { return servers[i].Addr.String() < servers[j].Addr.String() })5 days ago
continued
for _, s := range servers {
rtt := "n/a"
if s.AverageRTTSet {
rtt = s.AverageRTT.Round(time.Millisecond).String()
}
errStr := ""
if s.LastError != nil {
errStr = " lastError=" + s.LastError.Error()
}
log.Printf(" - %s kind=%s region=%s nodeType=%s az=%s rtt=%s%s",
s.Addr.String(),
s.Kind.String(),
orUnknown(tagValue(s.Tags, "region")),
orUnknown(tagValue(s.Tags, "nodeType")),
orUnknown(tagValue(s.Tags, "availabilityZone")),
rtt,
errStr,
)
}}
func (d *RegionDiag) reportRegions(topo description.Topology) {
seen := map[string]int{}
for _, s := range topo.Servers {
seen[tagValue(s.Tags, "region")]++
}
keys := make([]string, 0, len(seen))
for k := range seen {
keys = append(keys, k)
}
sort.Strings(keys)
log.Println("regions present in topology:")
for _, k := range keys {
log.Printf(" - %s (%d node(s))", orUnknown(k), seen[k])
}}
func (d *RegionDiag) reportTagEvaluation(topo description.Topology, rp *readpref.ReadPref) {
if rp == nil {
return
}
candidates := eligibleServers(topo)
log.Printf("tag evaluation against %d eligible candidate(s) for mode %s:", len(candidates), rp.Mode())
winner := -1
for i, ts := range rp.TagSets() {
var matched []string
for _, s := range candidates {
if tagSetMatches(ts, s.Tags) {
matched = append(matched, tagValue(s.Tags, "region"))
}
}
if len(matched) == 0 {
log.Printf(" tag set[%d] %s -> NO MATCH", i, formatTagSet(ts))
continue
}
sort.Strings(matched)
status := ""
if winner == -1 {
winner = i
status = " <== SELECTED (first match; later sets ignored)"
}
log.Printf(" tag set[%d] %s -> matches %d server(s) region(s)=%v%s",
i, formatTagSet(ts), len(matched), matched, status)
}
if winner < 0 {
log.Println(" RESULT: no tag set matched. Reads will fail with 'no server available'.")
return
}
winning := rp.TagSets()[winner]
if len(winning) == 0 {
log.Println(" RESULT: falling through to the EMPTY tag set. This is pure RTT-based")
log.Println(" 'nearest' selection, which is why you are seeing California reads.")
}}
type replSetMember struct {
Name string `bson:"name"`
StateStr string `bson:"stateStr"`
Health float64 `bson:"health"`}
type replSetStatus struct {
Members []replSetMember `bson:"members"`}
func (d *RegionDiag) reportReplicaSetStatus(ctx context.Context, client *mongo.Client) {
cmd := bson.D{{Key: "replSetGetStatus", Value: 1}}
var res replSetStatus
err := client.Database("admin").RunCommand(ctx, cmd).Decode(&res)
if err != nil {
log.Println("replSetGetStatus: unavailable (", err, ") - needs clusterMonitor role")
return
}
if len(res.Members) == 0 {
return
}
log.Println("replSetGetStatus members (server-side view, independent of driver):")
for _, m := range res.Members {
log.Printf(" - %s state=%s health=%.0f", m.Name, m.StateStr, m.Health)
}}
func eligibleServers(topo description.Topology) []description.Server {
var out []description.Server
for _, s := range topo.Servers {
switch s.Kind {
case description.RSPrimary, description.RSSecondary, description.RSMember:
out = append(out, s)
}
}
return out}
func tagSetMatches(ts tag.Set, serverTags tag.Set) bool {
for _, want := range ts {
found := false
for _, have := range serverTags {
if have.Name == want.Name && have.Value == want.Value {
found = true
break
}
}
if !found {
return false
}
}
return true}
func tagValue(ts tag.Set, name string) string {
for _, t := range ts {
if t.Name == name {
return t.Value
}
}
return ""}
func formatTagSet(ts tag.Set) string {
if len(ts) == 0 {
return "{} (empty - matches everything)"
}
parts := make([]string, 0, len(ts))
for _, t := range ts {
parts = append(parts, fmt.Sprintf("%s:%s", t.Name, t.Value))
}
return "{" + strings.Join(parts, ", ") + "}"}
func orUnknown(s string) string {
if s == "" {
return "<none>"
}
return s}
func replicaRegion() string {
r := os.Getenv("RAILWAY_REPLICA_REGION")
if r == "" {
return "<unset>"
}
return r}
func uriHost(uri string) string {
if uri == "" {
return "<unset>"
}
rest := uri
if i := strings.Index(rest, "@"); i >= 0 {
rest = rest[i+1:]
}
if i := strings.Index(rest, "/"); i >= 0 {
rest = rest[:i]
}
if i := strings.Index(rest, "?"); i >= 0 {
rest = rest[:i]
}
return rest}
5 days ago
main.go
package main
import (
"context"
"flag"
"log"
"net"
"os"
"time"
mongodrv "go.mongodb.org/mongo-driver/mongo"
"go.mongodb.org/mongo-driver/mongo/options"
"go.mongodb.org/mongo-driver/mongo/readpref"
"go.mongodb.org/mongo-driver/mongo/writeconcern"
"replica/internal/diag")
func main() {
region := flag.String("region", os.Getenv("RAILWAY_REPLICA_REGION"), "Railway region id to simulate (e.g. asia-southeast1-eqsg3a)")
uri := flag.String("uri", os.Getenv("DATABASE_URL"), "MongoDB connection string (defaults to $DATABASE_URL)")
flag.Parse()
if *uri == "" {
log.Fatal("no connection string: pass -uri or set DATABASE_URL")
}
log.Printf("using region %q (source: %s)", *region, regionSource(*region))
tagSets := diag.MapRegionToTags(*region)
rp, err := readpref.New(
readpref.NearestMode,
readpref.WithTagSets(tagSets...),
)
if err != nil {
log.Fatalf("failed to build read preference: %v", err)
}
d := diag.NewRegionDiag()
ctx, cancel := context.WithTimeout(context.Background(), 45*time.Second)
defer cancel()
dialer := &net.Dialer{KeepAlive: 300 * time.Second}
wc := writeconcern.New(writeconcern.W(1))
opts := options.Client().
ApplyURI(*uri).
SetDialer(dialer).
SetWriteConcern(wc).
SetReadPreference(rp).
SetServerMonitor(d.ServerMonitor()).
SetMonitor(d.CommandMonitor())
client, err := mongodrv.Connect(ctx, opts)
if err != nil {
log.Fatalf("connect: %v", err)
}
defer func() {
_ = client.Disconnect(context.Background())
}()
d.Report(ctx, client, rp)
d.SampleReads(ctx, client, 10)
d.StartRouteLogger(ctx, 5*time.Second)
time.Sleep(6 * time.Second)
d.LogRouteSummary()}
func regionSource(r string) string {
if os.Getenv("RAILWAY_REPLICA_REGION") == r {
return "RAILWAY_REPLICA_REGION"
}
return "-region flag"}