41
joinTimes []time.Time,
42
testInitialBootstrapFailure bool,
44
>
logger := log.NewTestLogger()
45
>
ctrl := gomock.NewController(t)
46
>
defer ctrl.Finish()
47
>
48
>
mockMgr := persistence.NewMockClusterMetadataManager(ctrl)
49
>
mockMgr.EXPECT().PruneClusterMembership(gomock.Any(), gomock.Any()).Return(nil).AnyTimes()
50
>
mockMgr.EXPECT().UpsertClusterMembership(gomock.Any(), gomock.Any()).Return(nil).AnyTimes()
51
>
52
>
cluster := &testCluster{
53
>
hostUUIDs: make([]string, size),
54
>
hostAddrs: make([]string, size),
55
>
hostInfoList: make([]membership.HostInfo, size),
56
>
rings: make([]*monitor, size),
57
>
channels: make([]*tchannel.Channel, size),
58
>
seedNode: seed,
59
>
}
60
>
61
>
for i := range size {
62
>
var err error
63
>
cluster.channels[i], err = tchannel.NewChannel(ringPopApp, nil)
64
>
if err != nil {
65
logger.Error("Failed to create tchannel", tag.Error(err))
66
return nil
67
}
69
>
err = cluster.channels[i].ListenAndServe(listenAddr)
70
>
if err != nil {
71
logger.Error("tchannel listen failed", tag.Error(err))
72
return nil
73
}
75
>
cluster.hostAddrs[i], err = buildBroadcastHostPort(cluster.channels[i].PeerInfo(), broadcastAddress)
76
>
if err != nil {
77
logger.Error("Failed to build broadcast hostport", tag.Error(err))
78
return nil
79
}
80
>
cluster.hostInfoList[i] = newHostInfo(cluster.hostAddrs[i], nil)
test_cluster.go
81
}
82
83
// if seed node is already supplied, use it; if not, set it
85
>
cluster.seedNode = cluster.hostAddrs[0]
86
>
}
87
>
logger.Info("seedNode", tag.Name(cluster.seedNode))
88
>
89
>
seedAddress, seedPort, err := splitHostPortTyped(cluster.seedNode)
90
>
if err != nil {
91
logger.Error("unable to split host port", tag.Error(err))
92
return nil
93
}
94
// MarshalBinary never fails for UUIDs
96
>
97
>
seedMember := &persistence.ClusterMember{
98
>
HostID: hostID,
99
>
RPCAddress: seedAddress,
100
>
RPCPort: seedPort,
101
>
SessionStart: time.Now().UTC(),
102
>
LastHeartbeat: time.Now().UTC(),
103
>
}
104
>
105
>
firstGetClusterMemberCall := testInitialBootstrapFailure
106
>
mockMgr.EXPECT().GetClusterMembers(gomock.Any(), gomock.Any()).DoAndReturn(
107
>
func(_ context.Context, _ *persistence.GetClusterMembersRequest) (*persistence.GetClusterMembersResponse, error) {
108
>
res := &persistence.GetClusterMembersResponse{ActiveMembers: []*persistence.ClusterMember{seedMember}}
109
>
110
>
hostID, _ := uuid.New().MarshalBinary()
111
>
if firstGetClusterMemberCall {
112
>
// The first time GetClusterMembers is invoked, we simulate returning a stale/bad heartbeat.
test_cluster.go
113
>
// All subsequent calls only return the single "good" seed member
114
>
// This ensures that we exercise the retry path in bootstrap properly.
115
>
badSeedMember := &persistence.ClusterMember{
116
>
HostID: hostID,
117
>
RPCAddress: seedAddress,
118
>
RPCPort: seedPort + 1,
119
>
SessionStart: time.Now().UTC(),
120
>
LastHeartbeat: time.Now().UTC(),
121
>
}
122
>
res = &persistence.GetClusterMembersResponse{ActiveMembers: []*persistence.ClusterMember{seedMember, badSeedMember}}
123
>
}
124
126
>
return res, nil
127
}).AnyTimes()
128
130
>
node := i
131
>
resolver := func() (string, error) {
132
>
return buildBroadcastHostPort(cluster.channels[node].PeerInfo(), broadcastAddress)
133
>
}
134
135
>
ringPop, err := ringpop.New(ringPopApp, ringpop.Channel(cluster.channels[i]), ringpop.AddressResolverFunc(resolver))
test_cluster.go
136
>
if err != nil {
137
logger.Error("failed to create ringpop instance", tag.Error(err))
138
return nil
139
}
140
>
_, port, _ := splitHostPortTyped(cluster.hostAddrs[i])
test_cluster.go
141
>
var joinTime time.Time
142
>
if i < len(joinTimes) {
143
joinTime = joinTimes[i]
144
}
146
>
serviceName,
147
>
config.ServicePortMap{serviceName: int(port)}, // use same port for "grpc" port
148
>
ringPop,
149
>
logger,
150
>
mockMgr,
151
>
resolver,
152
>
2*time.Second,
153
>
3*time.Second,
154
>
joinTime,
155
>
100,
156
>
)
157
>
cluster.rings[i].Start()
158
}
160
}
161