| 461 | for k := range m { |
| 462 | removed = append(removed, k) |
| 463 | } |
| 464 | return |
| 465 | } |
| 466 | |
| 467 | func (s *Daemon) runLinearizabilityTests(ctx context.Context) { |
| 468 | testers := map[string]linearizability.Tester{} |
| 469 | oldTargets := []string{} |
| 470 | ctx, cancel := context.WithCancel(ctx) |
| 471 | defer cancel() |
| 472 | |
| 473 | m := new(sync.Mutex) |
| 474 | |
| 475 | ticker := time.NewTicker(s.cfg.LinearizabilityTestInterval) |
| 476 | defer ticker.Stop() |
| 477 | for ; true; <-ticker.C { // first run without delay, then at interval |
| 478 | eg := new(errgroup.Group) |
| 479 | currentResults := map[string]linearizability.TestResult{} |
| 480 | targets, err := s.FetchGMs(s.cfg) |
| 481 | if err != nil { |
| 482 | log.Errorf("getting linearizability test targets: %v", err) |
| 483 | if len(oldTargets) > 0 { |
| 484 | linearizability.ProcessMonitoringResults(monPrefix, noTestResults(oldTargets), s.stats) |
| 485 | } else { |
| 486 | linearizability.ProcessMonitoringResults(monPrefix, noTestResults(defaultTargets), s.stats) |
| 487 | } |
| 488 | continue |
| 489 | } |
| 490 | log.Debugf("targets: %v, err: %v", targets, err) |
| 491 | // log when set of targets changes |
| 492 | added, removed := targetsDiff(oldTargets, targets) |
| 493 | if len(added) > 0 || len(removed) > 0 { |
| 494 | log.Infof("new set of linearizability test targets. Added: %v, Removed: %v. Resulting set: %v", added, removed, targets) |
| 495 | // we never remove testers, as stopping background listener goroutine is not easy |
| 496 | oldTargets = targets |
| 497 | } |
| 498 | |
| 499 | for _, server := range targets { |
| 500 | log.Debugf("talking to %s", server) |
| 501 | lt, found := testers[server] |
| 502 | if !found { |
| 503 | if s.cfg.SPTP { |
| 504 | lt, err = linearizability.NewSPTPHTTPTester(server, fmt.Sprintf("http://%s/", s.cfg.PTPClientAddress), s.cfg.LinearizabilityTestMaxGMOffset) |
| 505 | } else { |
| 506 | lt, err = linearizability.NewPTPTester(server, s.cfg.Iface, linearizability.IEEE1588) |
| 507 | } |
| 508 | if err != nil { |
| 509 | log.Errorf("creating tester: %v", err) |
| 510 | continue |
| 511 | } |
| 512 | testers[server] = lt |
| 513 | } |
| 514 | eg.Go(func() error { |
| 515 | res := lt.RunTest(ctx) |
| 516 | m.Lock() |
| 517 | currentResults[server] = res |
| 518 | m.Unlock() |
| 519 | return nil |
| 520 | }) |