234
ctx context.Context,
235
request *persistence.UpsertClusterMembershipRequest,
237
>
err := rpo.metadataManager.UpsertClusterMembership(ctx, request)
238
>
239
>
if err == nil {
240
>
hostID, err := uuid.FromBytes(request.HostID)
241
>
if err != nil {
242
return err
243
}
244
>
rpo.logger.Debug("Membership heartbeat upserted successfully",
monitor.go
245
>
tag.Address(request.RPCAddress.String()),
246
>
tag.Port(int(request.RPCPort)),
247
>
tag.HostID(hostID.String()))
248
}
249
251
}
252
253
// splitHostPortTyped expands upon net.SplitHostPort by providing type parsing.
254
>
func splitHostPortTyped(hostPort string) (net.IP, uint16, error) {
monitor.go
255
>
ipstr, portstr, err := net.SplitHostPort(hostPort)
256
>
if err != nil {
257
return nil, 0, err
258
}
259
260
>
broadcastAddress := net.ParseIP(ipstr)
monitor.go
261
>
broadcastPort, err := strconv.ParseUint(portstr, 10, 16)
262
>
if err != nil {
263
return nil, 0, err
264
}
265
266
>
return broadcastAddress, uint16(broadcastPort), nil
monitor.go
267
}
268
269
>
func (rpo *monitor) startHeartbeat(broadcastHostport string) error {
monitor.go
270
>
// Start by cleaning up expired records to avoid growth
271
>
err := rpo.metadataManager.PruneClusterMembership(rpo.lifecycleCtx, &persistence.PruneClusterMembershipRequest{MaxRecordsPruned: 10})
272
>
if err != nil {
273
return err
274
}
275
276
>
sessionStarted := time.Now().UTC()
monitor.go
277
>
278
>
// Parse and validate broadcast hostport
279
>
broadcastAddress, broadcastPort, err := splitHostPortTyped(broadcastHostport)
280
>
if err != nil {
281
return err
282
}
283
284
// Parse and validate existing service name
285
>
role, err := serviceNameToServiceTypeEnum(rpo.serviceName)
monitor.go
286
>
if err != nil {
287
return err
288
}
289
290
>
req := &persistence.UpsertClusterMembershipRequest{
monitor.go
291
>
Role: role,
292
>
RPCAddress: broadcastAddress,
293
>
RPCPort: broadcastPort,
294
>
SessionStart: sessionStarted,
295
>
RecordExpiry: upsertMembershipRecordExpiryDefault,
296
>
HostID: rpo.hostID,
297
>
}
298
>
299
>
// Upsert before fetching bootstrap hosts.
300
>
// This makes us discoverable by other Temporal cluster members
301
>
// Expire in 48 hours to allow for inspection of table by humans for debug scenarios.
302
>
// For bootstrapping, we filter to a much shorter duration on the
303
>
// read side by filtering on the last time a heartbeat was seen.
304
>
err = rpo.upsertMyMembership(rpo.lifecycleCtx, req)
305
>
if err == nil {
306
>
hostID, err := uuid.FromBytes(rpo.hostID)
307
>
if err != nil {
308
return err
309
}
310
>
rpo.logger.Info("Membership heartbeat upserted successfully",
monitor.go
311
>
tag.Address(broadcastAddress.String()),
312
>
tag.Port(int(broadcastPort)),
313
>
tag.HostID(hostID.String()))
314
>
315
>
rpo.startHeartbeatUpsertLoop(req)
316
}
317
319
}
320
321
>
func (rpo *monitor) fetchCurrentBootstrapHostports() ([]string, error) {
monitor.go
322
>
pageSize := 1000
323
>
set := make(map[string]struct{})
324
>
325
>
var nextPageToken []byte
326
>
327
>
for {
328
>
resp, err := rpo.metadataManager.GetClusterMembers(
329
>
rpo.lifecycleCtx,
330
>
&persistence.GetClusterMembersRequest{
331
>
LastHeartbeatWithin: healthyHostLastHeartbeatCutoff,
332
>
PageSize: pageSize,
333
>
NextPageToken: nextPageToken,
334
>
})
335
>
if err != nil {
336
return nil, err
337
}
338
339
// Dedupe on hostport
340
>
for _, host := range resp.ActiveMembers {
monitor.go
341
>
set[net.JoinHostPort(host.RPCAddress.String(), convert.Uint16ToString(host.RPCPort))] = struct{}{}
342
>
}
343
>
nextPageToken = resp.NextPageToken
344
>
345
>
// Stop iterating once we have either 500 unique ip:port combos or there is no more results.
346
>
if len(nextPageToken) == 0 || len(set) >= 500 {
347
>
bootstrapHostPorts := make([]string, 0, len(set))
348
>
for k := range set {
349
>
bootstrapHostPorts = append(bootstrapHostPorts, k)
350
>
}
351
352
>
rpo.logger.Info("bootstrap hosts fetched", tag.BootstrapHostPorts(strings.Join(bootstrapHostPorts, ",")))
monitor.go
353
>
return bootstrapHostPorts, nil
354
}
355
}
356
}
357
358
>
func (rpo *monitor) startHeartbeatUpsertLoop(request *persistence.UpsertClusterMembershipRequest) {
monitor.go
359
>
loopUpsertMembership := func() {
360
>
for {
361
>
select {
362
case <-rpo.lifecycleCtx.Done():
363
return
365
}
366
>
err := rpo.upsertMyMembership(rpo.lifecycleCtx, request)
monitor.go
367
>
if err != nil {
368
rpo.logger.Error("Membership upsert failed.", tag.Error(err))
369
}
370
371
>
jitter := math.Round(rand.Float64() * 5)
monitor.go
372
>
time.Sleep(time.Second * time.Duration(10+jitter))
373
}
374
}
375
377
}
378