134
}
135
136
>
func (s *ClusterMetadataManagerSuite) validateUpsert(req *p.UpsertClusterMembershipRequest, resp *p.GetClusterMembersResponse, err error) {
cluster_metadata_manager.go
137
>
s.Nil(err)
138
>
s.NotNil(resp)
139
>
s.NotEmpty(resp.ActiveMembers)
140
>
s.Equal(len(resp.ActiveMembers), 1)
141
>
// Have to round to 1 second due to SQL implementations. Cassandra truncates at 1ms.
142
>
s.Equal(resp.ActiveMembers[0].SessionStart.Round(time.Second), req.SessionStart.Round(time.Second))
143
>
s.Equal(resp.ActiveMembers[0].RPCAddress.String(), req.RPCAddress.String())
144
>
s.Equal(resp.ActiveMembers[0].RPCPort, req.RPCPort)
145
>
s.True(resp.ActiveMembers[0].RecordExpiry.After(time.Now().UTC()))
146
>
s.Equal(resp.ActiveMembers[0].HostID, req.HostID)
147
>
s.Equal(resp.ActiveMembers[0].Role, req.Role)
148
>
}
149
150
// TestClusterMembershipReadFiltersCorrectly verifies that we can UpsertClusterMembership and read our result using filters
152
>
now := time.Now().UTC()
153
>
hostID, err := uuid.New().MarshalBinary()
154
>
s.NoError(err)
155
>
req := &p.UpsertClusterMembershipRequest{
156
>
HostID: hostID,
157
>
RPCAddress: net.ParseIP("127.0.0.2"),
158
>
RPCPort: 123,
159
>
Role: p.Frontend,
160
>
SessionStart: now,
161
>
RecordExpiry: time.Second * 4,
162
>
}
163
>
164
>
err = s.ClusterMetadataManager.UpsertClusterMembership(s.ctx, req)
165
>
s.Nil(err)
166
>
167
>
resp, err := s.ClusterMetadataManager.GetClusterMembers(
168
>
s.ctx,
169
>
&p.GetClusterMembersRequest{LastHeartbeatWithin: time.Minute * 10, HostIDEquals: req.HostID},
170
>
)
171
>
172
>
s.validateUpsert(req, resp, err)
173
>
174
>
time.Sleep(time.Second * 1)
175
>
resp, err = s.ClusterMetadataManager.GetClusterMembers(
176
>
s.ctx,
177
>
&p.GetClusterMembersRequest{LastHeartbeatWithin: time.Millisecond, HostIDEquals: req.HostID},
178
>
)
179
>
180
>
s.Nil(err)
181
>
s.NotNil(resp)
182
>
s.Empty(resp.ActiveMembers)
183
>
184
>
resp, err = s.ClusterMetadataManager.GetClusterMembers(
185
>
s.ctx,
186
>
&p.GetClusterMembersRequest{RoleEquals: p.Matching},
187
>
)
188
>
189
>
s.Nil(err)
190
>
s.NotNil(resp)
191
>
s.Empty(resp.ActiveMembers)
192
>
193
>
resp, err = s.ClusterMetadataManager.GetClusterMembers(
194
>
s.ctx,
195
>
&p.GetClusterMembersRequest{SessionStartedAfter: time.Now().UTC()},
196
>
)
197
>
198
>
s.Nil(err)
199
>
s.NotNil(resp)
200
>
s.Empty(resp.ActiveMembers)
201
>
202
>
resp, err = s.ClusterMetadataManager.GetClusterMembers(
203
>
s.ctx,
204
>
&p.GetClusterMembersRequest{SessionStartedAfter: now.Add(-time.Minute), RPCAddressEquals: req.RPCAddress, HostIDEquals: req.HostID},
205
>
)
206
>
207
>
s.validateUpsert(req, resp, err)
208
>
s.waitForPrune(5 * time.Second)
209
>
}
210
211
// TestClusterMembershipUpsertExpiresCorrectly verifies RecordExpiry functions properly for ClusterMembership records