monitor.go ×40

Frontier kind: Code frontier

unlabeled · c_2b67f95dfd6f

13 tests · 2291 LOC · 86 files · introduces 0 tests · 336 LOC · 5 files

Introduces — evidence that enters the hierarchy at this concept

Code
72 ranges336 lines · 5 files
Tests
0 tests

Contains — complete concept membership

All code (extent)
358 ranges2291 lines · 86 files · Browse complete extent
All tests (intent)
13 testsBrowse complete intent

Neighbourhood graph

The orange circle is the focus. Violet and green circles are every ancestor and descendant, broader and narrower, at any distance; blue squares and pink diamonds are the introduced files and exact introduced tests of every visible concept, not only the focus's. Arrows point from broader to narrower concepts and bridge only concepts omitted from this view. Undirected links show source or test introduction. Concept and file size follows LOC; exact test nodes use test-count units.

Introduced files, introduced tests, and structurally relevant concept specialization

In the embedded map, ordinary wheel input scrolls the page; use the visible controls to zoom and drag to pan. Open the full-screen map for canvas navigation: wheel pans, Ctrl/Command plus wheel zooms, and arrow keys pan when this region is focused. On touch screens, open the full-screen map to pan or pinch. If JavaScript or WebGL is unavailable, use the native relationship evidence on this page.

Graph controls are ready.

Interactive rendering requires JavaScript and WebGL. Use the native relationship evidence on this page while the interactive map is unavailable.

Native relationship evidence

Every exact file and test below is linked only from the concept that introduces it.

Introduced tests

Every collected test enters the hierarchy at exactly one concept.

No tests are introduced at this concept. Its intent tests are introduced by other concepts.

Introduced code

Every collected source range enters the hierarchy at exactly one concept.

5 files ranked by introduced lines: 336 introduced LOC across 72 ranges. Expand a file to inspect source; the > gutter marks introduced lines.

go.temporal.io/server/common/membership/ringpop/monitor.go 178 introduced LOC · 40 ranges

Open complete file

84 joinTime time.Time,
85 replicaPoints int,
86 > ) *monitor { monitor.go
87 > lifecycleCtx, lifecycleCancel := context.WithCancel(context.Background())
88 > lifecycleCtx = headers.SetCallerInfo(
89 > lifecycleCtx,
90 > headers.SystemBackgroundHighCallerInfo,
91 > )
92 > hostID, _ := uuid.New().MarshalBinary()
93 > // MarshalBinary should never error.
94 >
95 > rpo := &monitor{
96 > status: common.DaemonStatusInitialized,
97 >
98 > lifecycleCtx: lifecycleCtx,
99 > lifecycleCancel: lifecycleCancel,
100 >
101 > serviceName: serviceName,
102 > services: services,
103 > rp: rp,
104 > rings: make(map[primitives.ServiceName]*serviceResolver),
105 > replicaPoints: replicaPoints,
106 > logger: logger,
107 > metadataManager: metadataManager,
108 > broadcastHostPortResolver: broadcastHostPortResolver,
109 > hostID: hostID,
110 > initialized: future.NewFuture[struct{}](),
111 > maxJoinDuration: maxJoinDuration,
112 > propagationTime: propagationTime,
113 > joinTime: joinTime,
114 > }
115 > for service, port := range services {
116 > rpo.rings[service] = newServiceResolver(service, port, rp, replicaPoints, logger)
117 > }
118 > return rpo
119 }
120
122 // it's safe for Stop() to run, which is at any point when we are neither updating the status field nor
123 // starting rings
124 > func (rpo *monitor) Start() { monitor.go
125 > rpo.stateLock.Lock()
126 > if rpo.status != common.DaemonStatusInitialized {
127 rpo.stateLock.Unlock()
128 return
129 }
130 > rpo.status = common.DaemonStatusStarted monitor.go
131 > rpo.stateLock.Unlock()
132 >
133 > broadcastAddress, err := rpo.broadcastHostPortResolver()
134 > if err != nil {
135 rpo.logger.Fatal("unable to resolve broadcast address", tag.Error(err))
136 }
140 // we must know our seed nodes before bootstrapping
141
142 > if err = rpo.startHeartbeat(broadcastAddress); err != nil { monitor.go
143 rpo.logger.Fatal("unable to initialize membership heartbeats", tag.Error(err))
144 }
145
146 > if err = rpo.bootstrapRingPop(); err != nil { monitor.go
147 // Stop() called during Start()'s execution. This is ok
148 if strings.Contains(err.Error(), "destroyed while attempting to join") {
152 }
153
154 > labels, err := rpo.rp.Labels() monitor.go
155 > if err != nil {
156 rpo.logger.Fatal("unable to get ringpop labels", tag.Error(err))
157 }
158
159 > if until := time.Until(rpo.joinTime); until > 0 && until.Seconds() < maxScheduledEventTimeSeconds { monitor.go
160 if err = labels.Set(startAtKey, strconv.FormatInt(rpo.joinTime.Unix(), 10)); err != nil {
161 rpo.logger.Fatal("unable to set ringpop label", tag.Error(err), tag.Key(startAtKey))
168 }
169
170 > if err = labels.Set(portKey, strconv.Itoa(rpo.services[rpo.serviceName])); err != nil { monitor.go
171 rpo.logger.Fatal("unable to set ringpop label", tag.Error(err), tag.Key(portKey))
172 }
173
174 // This label should be set last, it's used as the prediciate for finding members for rings.
175 > if err = labels.Set(roleKey, string(rpo.serviceName)); err != nil { monitor.go
176 rpo.logger.Fatal("unable to set ringpop label", tag.Error(err), tag.Key(roleKey))
177 }
178
179 // Our individual rings may not support concurrent start/stop calls so we reacquire the state lock while acting upon them
180 > rpo.stateLock.Lock() monitor.go
181 > for _, ring := range rpo.rings {
182 > ring.Start()
183 > }
184 > rpo.stateLock.Unlock()
185 >
186 > rpo.initialized.Set(struct{}{}, nil)
187 }
188
189 // bootstrap ring pop service by discovering the bootstrap hosts and joining the ring pop cluster
190 > func (rpo *monitor) bootstrapRingPop() error { monitor.go
191 > policy := backoff.NewExponentialRetryPolicy(healthyHostLastHeartbeatCutoff / 2).
192 > WithBackoffCoefficient(1).
193 > WithMaximumAttempts(maxBootstrapRetries)
194 > op := func() error {
195 > hostPorts, err := rpo.fetchCurrentBootstrapHostports()
196 > if err != nil {
197 return err
198 }
199
200 > bootParams := &swim.BootstrapOptions{ monitor.go
201 > ParallelismFactor: 10,
202 > JoinSize: 1,
203 > MaxJoinDuration: rpo.maxJoinDuration,
204 > DiscoverProvider: statichosts.New(hostPorts...),
205 > }
206 >
207 > _, err = rpo.rp.Bootstrap(bootParams)
208 > if err != nil {
209 rpo.logger.Warn("unable to bootstrap ringpop. retrying", tag.Error(err))
210 }
211 > return err monitor.go
212 }
213
214 > if err := backoff.ThrottleRetry(op, policy, nil); err != nil { monitor.go
215 return fmt.Errorf("exhausted all retries: %w", err)
216 }
217 > return nil monitor.go
218 }
219
234 ctx context.Context,
235 request *persistence.UpsertClusterMembershipRequest,
236 > ) error { monitor.go
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
250 > return err monitor.go
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
318 > return err monitor.go
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
364 > default: monitor.go
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
376 > go loopUpsertMembership() monitor.go
377 }
378
425 }
426
427 > func (rpo *monitor) GetResolver(service primitives.ServiceName) (membership.ServiceResolver, error) { monitor.go
428 > ring, found := rpo.rings[service]
429 > if !found {
430 return nil, membership.ErrUnknownService
431 }
432 > return ring, nil monitor.go
433 }
434
450 }
451
452 > func replaceServicePort(address string, servicePort int) (string, error) { monitor.go
453 > host, _, err := net.SplitHostPort(address)
454 > if err != nil {
455 return "", membership.ErrIncorrectAddressFormat
456 }
457 > return net.JoinHostPort(host, convert.IntToString(servicePort)), nil monitor.go
458 }
459
476 }
477
478 > func serviceNameToServiceTypeEnum(name primitives.ServiceName) (persistence.ServiceType, error) { monitor.go
479 > if serviceType, ok := serviceNameToServiceTypeEnumMap[name]; ok {
480 > return serviceType, nil
481 > }
482
483 return persistence.All, fmt.Errorf("unable to parse servicename '%s'", name)
go.temporal.io/server/common/membership/ringpop/service_resolver.go 144 introduced LOC · 27 ranges

Open complete file

99 replicaPoints int,
100 logger log.Logger,
101 > ) *serviceResolver { service_resolver.go
102 > resolver := &serviceResolver{
103 > service: service,
104 > port: port,
105 > rp: rp,
106 > replicaPoints: replicaPoints,
107 > refreshChan: make(chan struct{}),
108 > shutdownCh: make(chan struct{}),
109 > logger: log.With(logger, tag.ComponentServiceResolver, tag.Service(service)),
110 > scheduledRefreshMap: make(map[int64]*time.Timer),
111 > listeners: make(map[string]chan<- *membership.ChangedEvent),
112 > }
113 > resolver.ringAndHosts.Store(ringAndHosts{
114 > ring: newHashRing(replicaPoints),
115 > hosts: make(map[string]*hostInfo),
116 > })
117 > return resolver
118 > }
119
120 > func newHashRing(replicaPoints int) *hashring.HashRing { service_resolver.go
121 > return hashring.New(farm.Fingerprint32, replicaPoints)
122 > }
123
124 // Start starts the oracle
125 > func (r *serviceResolver) Start() { service_resolver.go
126 > r.rp.AddListener(r)
127 > if err := r.refresh(refreshModeAlways); err != nil {
128 r.logger.Fatal("unable to start ring pop service resolver", tag.Error(err))
129 }
130
131 > r.shutdownWG.Add(1) service_resolver.go
132 > go r.refreshRingWorker()
133 }
134
247 func (r *serviceResolver) HandleEvent(
248 event events.Event,
250 > // We only about membership.ChangeEvent. Normally ringpop converts membership.ChangeEvent
251 > // into events.RingChangedEvent when its internal hash ring changes, but since we construct
252 > // our own hash rings with filtering, we have to handle the lower-level event ourselves.
253 > if _, ok := event.(rpmembership.ChangeEvent); ok {
254 > r.logger.Debug("Received a ring changed event")
255 > // Note that we receive events asynchronously, possibly out of order.
256 > // We cannot rely on the content of the event, rather we load everything
257 > // from ringpop when we get a notification that something changed.
258 > if err := r.refresh(refreshModeAlways); err != nil {
259 r.logger.Error("error refreshing ring when receiving a ring changed event", tag.Error(err))
260 }
262 }
263
264 > func (r *serviceResolver) refresh(mode refreshMode) error { service_resolver.go
265 > var event *membership.ChangedEvent
266 > var err error
267 > defer func() {
268 > if event != nil {
269 > r.emitEvent(event)
270 > }
271 }()
272
273 > r.refreshLock.Lock() service_resolver.go
274 > defer r.refreshLock.Unlock()
275 >
276 > if mode == refreshModeLazy && r.lastRefreshTime.After(time.Now().UTC().Add(-minRefreshInternal)) {
277 return nil // refreshed too recently
278 }
279
280 > event, err = r.refreshLocked() service_resolver.go
281 > return err
282 }
283
284 > func (r *serviceResolver) refreshLocked() (*membership.ChangedEvent, error) { service_resolver.go
285 > hosts, nextEvent, err := r.getReachableMembers()
286 > if err != nil {
287 return nil, err
288 }
289
290 // if we found an add/remove event, schedule another refresh right at that time
291 > r.scheduleRefresh(nextEvent) service_resolver.go
292 >
293 > newMembersMap, changedEvent := r.compareMembers(hosts)
294 > if changedEvent == nil {
295 > return nil, nil
296 > }
297
298 > ring := newHashRing(r.replicaPoints) service_resolver.go
299 > ring.AddMembers(util.MapSlice(hosts, func(h *hostInfo) rpmembership.Member { return h })...)
300
301 > r.lastRefreshTime = time.Now().UTC() service_resolver.go
302 > r.ringAndHosts.Store(ringAndHosts{
303 > ring: ring,
304 > hosts: newMembersMap,
305 > })
306 >
307 > addrs := util.MapSlice(hosts, func(h *hostInfo) string { return h.summary() })
308 > slices.Sort(addrs)
309 > r.logger.Info("Current reachable members", tag.Addresses(addrs))
310 >
311 > return changedEvent, nil
312 }
313
314 > func (r *serviceResolver) scheduleRefresh(nextEvent int64) { service_resolver.go
315 > if nextEvent == 0 {
316 > return
317 > }
318 if _, ok := r.scheduledRefreshMap[nextEvent]; ok {
319 return // already have a timer scheduled for this time
333 }
334
335 > func (r *serviceResolver) getReachableMembers() ([]*hostInfo, int64, error) { service_resolver.go
336 > members, err := r.rp.GetReachableMemberObjects(swim.MemberWithLabelAndValue(roleKey, string(r.service)))
337 > if err != nil {
338 return nil, 0, err
339 }
342 // need to keep track of one event since we'll refresh at that time and find the next one.
343 // Note that nextEvent is mutated by the filter functions below.
344 > nowUnix := time.Now().Unix() service_resolver.go
345 > nextEvent := int64(math.MaxInt64)
346 >
347 > // Filter by startAt
348 > members = slices.DeleteFunc(members, func(member swim.Member) bool {
349 > startAt, err := parseIntLabel(member, startAtKey)
350 > if err != nil {
351 > return false // ignore label if missing or can't parse
352 > } else if startAt <= nowUnix {
353 return false // start time is in the past
354 }
359
360 // Filter by stopAt
361 > members = slices.DeleteFunc(members, func(member swim.Member) bool { service_resolver.go
362 > stopAt, err := parseIntLabel(member, stopAtKey)
363 > if err != nil {
364 > return false // ignore label if missing or can't parse
365 > } else if stopAt > nowUnix {
366 // stop time is in the future: schedule refresh at that time
367 nextEvent = min(nextEvent, stopAt)
372
373 // Turn swim.Members into hostInfo
374 > hosts := make([]*hostInfo, len(members)) service_resolver.go
375 > for i, member := range members {
376 > servicePort := r.port
377 >
378 > // Each temporal service in the ring should advertise which port it has its gRPC listener
379 > // on via a service label. If we cannot find the label, we will assume that the
380 > // temporal service is listening on the same port that this node is listening on.
381 > servicePortLabel, ok := member.Label(portKey)
382 > if ok {
383 > servicePort, err = strconv.Atoi(servicePortLabel)
384 > if err != nil {
385 return nil, 0, err
386 }
389 }
390
391 > hostPort, err := replaceServicePort(member.Address, servicePort) service_resolver.go
392 > if err != nil {
393 return nil, 0, err
394 }
395
396 // We can share member.Labels without copying since we never modify it.
397 > hosts[i] = newHostInfo(hostPort, member.Labels) service_resolver.go
398 }
399
400 > if nextEvent == math.MaxInt64 { service_resolver.go
401 > nextEvent = 0
402 > }
403 > return hosts, nextEvent, nil
404 }
405
406 > func (r *serviceResolver) emitEvent(event *membership.ChangedEvent) { service_resolver.go
407 > // Notify listeners
408 > r.listenerLock.RLock()
409 > defer r.listenerLock.RUnlock()
410 >
411 > for name, ch := range r.listeners {
412 select {
413 case ch <- event:
418 }
419
420 > func (r *serviceResolver) refreshRingWorker() { service_resolver.go
421 > defer r.shutdownWG.Done()
422 >
423 > refreshTicker := time.NewTicker(defaultRefreshInterval)
424 > defer refreshTicker.Stop()
425 >
426 > for {
427 > select {
428 case <-r.shutdownCh:
429 return
474 // buildBroadcastHostPort return the listener hostport from an existing tchannel
475 // and overrides the address with broadcastAddress if specified
476 > func buildBroadcastHostPort(listenerPeerInfo tchannel.LocalPeerInfo, broadcastAddress string) (string, error) { service_resolver.go
477 > // Ephemeral port check copied from ringpop-go/ringpop.go/channelAddressResolver
478 > // Check that TChannel is listening on a real hostport. By default,
479 > // TChannel listens on an ephemeral host/port. The real port is then
480 > // assigned by the OS when ListenAndServe is called. If the hostport is
481 > // ephemeral, it means TChannel is not yet listening and the hostport
482 > // cannot be resolved.
483 > if listenerPeerInfo.IsEphemeralHostPort() {
484 return "", ringpop.ErrEphemeralAddress
485 }
486
487 // Parse listener hostport
488 > listenerIPString, port, err := net.SplitHostPort(listenerPeerInfo.HostPort) service_resolver.go
489 > if err != nil {
490 return "", err
491 }
492
493 // Broadcast IP override
494 > if broadcastAddress != "" { service_resolver.go
495 > // Parse supplied broadcastAddress override
496 > ip := net.ParseIP(broadcastAddress)
497 > if ip == nil {
498 return "", errors.New("broadcastAddress set but unknown failure encountered while parsing")
499 }
500
501 // If no errors, use the parsed IP with the port from our listener
502 > return net.JoinHostPort(ip.String(), port), nil service_resolver.go
503 }
504
516
517 // parseIntLabel returns the value of the given label as an integer.
518 > func parseIntLabel(member swim.Member, label string) (int64, error) { service_resolver.go
519 > str, ok := member.Label(label)
520 > if !ok {
521 > return 0, errMissingLabel
522 > }
523 return strconv.ParseInt(str, 10, 64)
524 }
go.temporal.io/server/common/log/tag/tags.go 6 introduced LOC · 2 ranges

Open complete file

424
425 // Addresses returns tag for Addresses
426 > func Addresses(ads []string) ZapTag { tags.go
427 > return NewStringsTag("addresses", ads)
428 > }
429
430 // ListenerName returns tag for ListenerName
954
955 // BootstrapHostPorts returns tag for bootstrap host ports
956 > func BootstrapHostPorts(s string) ZapTag { tags.go
957 > return NewStringTag("bootstrap-hostports", s)
958 > }
959
960 // TLSCertFile returns tag for TLS cert file name
go.temporal.io/server/common/membership/ringpop/hostinfo.go 5 introduced LOC · 2 ranges

Open complete file

37
38 // Identity implements ringpop's Membership interface
39 > func (hi *hostInfo) Identity() string { hostinfo.go
40 > // For now, we just use the address as the identity.
41 > return hi.addr
42 > }
43
44 // Label implements ringpop's Membership interface
54 for k, v := range hi.labels {
55 switch k {
56 > case roleKey, portKey: hostinfo.go
57 // skip these, they can be determined from context
58 default:
go.temporal.io/server/common/convert/convert.go 3 introduced LOC · 1 range

Open complete file

32 }
33
34 > func Uint16ToString(v uint16) string { convert.go
35 > return strconv.FormatUint(uint64(v), 10)
36 > }
37
38 func Int64SetToSlice(