| 1 | // SPDX-License-Identifier: GPL-3.0-or-later |
| 2 | |
| 3 | package chrony |
| 4 | |
| 5 | import ( |
| 6 | "bufio" |
| 7 | "bytes" |
| 8 | "errors" |
| 9 | "fmt" |
| 10 | "strconv" |
| 11 | "strings" |
| 12 | "time" |
| 13 | |
| 14 | "github.com/netdata/netdata/go/plugins/plugin/go.d/pkg/oldmetrix" |
| 15 | ) |
| 16 | |
| 17 | const scaleFactor = 1000000000 |
| 18 | |
| 19 | const ( |
| 20 | // https://github.com/mlichvar/chrony/blob/7daf34675a5a2487895c74d1578241ca91a4eb70/ntp.h#L70-L75 |
| 21 | leapStatusNormal = 0 |
| 22 | leapStatusInsertSecond = 1 |
| 23 | leapStatusDeleteSecond = 2 |
| 24 | leapStatusUnsynchronised = 3 |
| 25 | ) |
| 26 | |
| 27 | func (c *Collector) collect() (map[string]int64, error) { |
| 28 | if c.conn == nil { |
| 29 | client, err := c.newConn(c.Config) |
| 30 | if err != nil { |
| 31 | return nil, err |
| 32 | } |
| 33 | c.conn = client |
| 34 | } |
| 35 | |
| 36 | mx := make(map[string]int64) |
| 37 | |
| 38 | if err := c.collectTracking(mx); err != nil { |
| 39 | return nil, err |
| 40 | } |
| 41 | if err := c.collectActivity(mx); err != nil { |
| 42 | return mx, err |
| 43 | } |
| 44 | if c.exec != nil { |
| 45 | if err := c.collectServerStats(mx); err != nil { |
| 46 | c.Warning(err) |
| 47 | c.exec = nil |
| 48 | } else { |
| 49 | c.addServerStatsChartsOnce.Do(c.addServerStatsCharts) |
| 50 | } |
| 51 | } |
| 52 | |
| 53 | return mx, nil |
| 54 | } |
| 55 | |
| 56 | func (c *Collector) collectTracking(mx map[string]int64) error { |
| 57 | reply, err := c.conn.tracking() |
| 58 | if err != nil { |
| 59 | return fmt.Errorf("error on collecting tracking: %v", err) |
| 60 | } |
| 61 | |
| 62 | mx["stratum"] = int64(reply.Stratum) |
| 63 | mx["leap_status_normal"] = oldmetrix.Bool(reply.LeapStatus == leapStatusNormal) |
| 64 | mx["leap_status_insert_second"] = oldmetrix.Bool(reply.LeapStatus == leapStatusInsertSecond) |
| 65 | mx["leap_status_delete_second"] = oldmetrix.Bool(reply.LeapStatus == leapStatusDeleteSecond) |
| 66 | mx["leap_status_unsynchronised"] = oldmetrix.Bool(reply.LeapStatus == leapStatusUnsynchronised) |
| 67 | mx["root_delay"] = int64(reply.RootDelay * scaleFactor) |
| 68 | mx["root_dispersion"] = int64(reply.RootDispersion * scaleFactor) |
| 69 | mx["skew"] = int64(reply.SkewPPM * scaleFactor) |
| 70 | mx["last_offset"] = int64(reply.LastOffset * scaleFactor) |
| 71 | mx["rms_offset"] = int64(reply.RMSOffset * scaleFactor) |
| 72 | mx["update_interval"] = int64(reply.LastUpdateInterval * scaleFactor) |
| 73 | // handle chrony restarts |
| 74 | if reply.RefTime.Year() != 1970 { |
| 75 | mx["ref_measurement_time"] = time.Now().Unix() - reply.RefTime.Unix() |
| 76 | } |
| 77 | mx["residual_frequency"] = int64(reply.ResidFreqPPM * scaleFactor) |
| 78 | // https://github.com/mlichvar/chrony/blob/5b04f3ca902e5d10aa5948fb7587d30b43941049/client.c#L1706 |
| 79 | mx["current_correction"] = abs(int64(reply.CurrentCorrection * scaleFactor)) |
| 80 | mx["frequency"] = abs(int64(reply.FreqPPM * scaleFactor)) |
| 81 | |
| 82 | return nil |
| 83 | } |
| 84 | |
| 85 | func (c *Collector) collectActivity(mx map[string]int64) error { |
| 86 | reply, err := c.conn.activity() |
| 87 | if err != nil { |
| 88 | return fmt.Errorf("error on collecting activity: %v", err) |
| 89 | } |
| 90 | |
| 91 | mx["online_sources"] = int64(reply.Online) |
| 92 | mx["offline_sources"] = int64(reply.Offline) |
| 93 | mx["burst_online_sources"] = int64(reply.BurstOnline) |
| 94 | mx["burst_offline_sources"] = int64(reply.BurstOffline) |
| 95 | mx["unresolved_sources"] = int64(reply.Unresolved) |
| 96 | |
| 97 | return nil |
| 98 | } |
| 99 | |
| 100 | func (c *Collector) collectServerStats(mx map[string]int64) error { |
| 101 | bs, err := c.exec.serverStats() |
| 102 | if err != nil { |
| 103 | return fmt.Errorf("error on collecting server stats: %v", err) |
| 104 | } |
| 105 | |
| 106 | sc := bufio.NewScanner(bytes.NewReader(bs)) |
| 107 | var n int |
| 108 | |
| 109 | for sc.Scan() { |
| 110 | key, value, ok := strings.Cut(sc.Text(), ":") |
| 111 | if !ok { |
| 112 | continue |
| 113 | } |
| 114 | |
| 115 | key, value = strings.TrimSpace(key), strings.TrimSpace(value) |
| 116 | |
| 117 | switch key { |
| 118 | case "NTP packets received", |
| 119 | "NTP packets dropped", |
| 120 | "Command packets received", |
| 121 | "Command packets dropped": |
| 122 | if v, err := strconv.ParseInt(value, 10, 64); err == nil { |
| 123 | key = strings.ToLower(strings.ReplaceAll(key, " ", "_")) |
| 124 | mx[key] = v |
| 125 | n++ |
| 126 | } |
| 127 | } |
| 128 | } |
| 129 | |
| 130 | if n == 0 { |
| 131 | return errors.New("no server stats metrics found in the response") |
| 132 | } |
| 133 | |
| 134 | return nil |
| 135 | } |
| 136 | |
| 137 | func abs(v int64) int64 { |
| 138 | if v < 0 { |
| 139 | return -v |
| 140 | } |
| 141 | return v |
| 142 | } |