master
go 214 lines 5.54 KB
Raw
1 // SPDX-License-Identifier: GPL-3.0-or-later
2
3 package cato_networks
4
5 import (
6 "context"
7 "errors"
8 "fmt"
9 "strings"
10 "time"
11 )
12
13 type discoveryState struct {
14 siteIDs []string
15 siteNames map[string]string
16 fetchedAt time.Time
17 skippedBySelector int
18 }
19
20 func (c *Collector) collect(ctx context.Context) error {
21 if c.client == nil {
22 return errors.New("Cato client is not initialized")
23 }
24
25 if err := c.refreshDiscovery(ctx, false); err != nil {
26 return wrapCatoOperationError("site discovery", err)
27 }
28 if len(c.discovery.siteIDs) == 0 {
29 return errors.New("no Cato sites discovered")
30 }
31
32 sites, order, err := c.collectSnapshot(ctx)
33 if err != nil {
34 return wrapCatoOperationError("account snapshot", err)
35 }
36 if len(sites) == 0 {
37 return errors.New("no Cato sites returned by account snapshot")
38 }
39 c.pruneUnselectedSites(sites, &order)
40
41 if err := c.collectMetrics(ctx, sites); err != nil {
42 errorClass := classifyCatoError(err)
43 c.warnRecoverable(warningKeyMetrics, errorClass, "account metrics collection incomplete, error_class=%s", errorClass)
44 }
45
46 if err := c.collectBGP(ctx, sites, order); err != nil {
47 errorClass := classifyCatoError(err)
48 c.warnRecoverable(warningKeyBGP, errorClass, "BGP status collection incomplete, error_class=%s", errorClass)
49 }
50
51 c.pruneUnselectedSites(sites, &order)
52
53 now := c.now()
54 topo, err := buildTopology(c.AccountID, sites, order, now)
55 if err != nil {
56 return fmt.Errorf("build topology: %w", err)
57 }
58 c.topology.Publish(topo)
59
60 c.writeMetrics(sites, order)
61
62 return nil
63 }
64
65 func (c *Collector) refreshDiscovery(ctx context.Context, force bool) error {
66 now := c.now()
67 if !force && len(c.discovery.siteIDs) > 0 {
68 if now.Sub(c.discovery.fetchedAt) < seconds(defaultDiscoveryEvery) {
69 return nil
70 }
71 }
72
73 limit := int64(defaultDiscoveryLimit)
74 var from int64
75 siteNames := make(map[string]string)
76 var siteIDs []string
77
78 for {
79 if from/limit >= maxDiscoveryPages {
80 err := fmt.Errorf("entityLookup pagination exceeded %d pages", maxDiscoveryPages)
81 if c.useCachedDiscoveryAfterRefreshFailure(now, force, err) {
82 return nil
83 }
84 return err
85 }
86
87 res, err := c.client.LookupSites(ctx, c.AccountID, limit, from)
88 if err != nil {
89 if c.useCachedDiscoveryAfterRefreshFailure(now, force, err) {
90 return nil
91 }
92 return err
93 }
94
95 items := res.GetEntityLookup().GetItems()
96 for _, item := range items {
97 entity := item.GetEntity()
98 siteID := strings.TrimSpace(entity.GetID())
99 if siteID == "" {
100 continue
101 }
102 siteIDs = append(siteIDs, siteID)
103 if name := derefZero(entity.GetName()); name != "" {
104 siteNames[siteID] = name
105 }
106 }
107
108 total := derefZero(res.GetEntityLookup().GetTotal())
109 from += int64(len(items))
110 if len(items) == 0 || (total > 0 && from >= total) || int64(len(items)) < limit {
111 break
112 }
113 }
114
115 selectedSiteIDs, skippedBySelector := c.selectSites(siteIDs, siteNames)
116 c.discovery = discoveryState{
117 siteIDs: selectedSiteIDs,
118 siteNames: siteNames,
119 fetchedAt: now,
120 skippedBySelector: skippedBySelector,
121 }
122 if skippedBySelector > 0 {
123 c.Limit("cato:site_selector", 1, recurringLogEvery).
124 Infof("selected %d of %d Cato site(s) after applying site_selector", len(selectedSiteIDs), len(siteIDs))
125 }
126 return nil
127 }
128
129 func (c *Collector) selectSites(siteIDs []string, siteNames map[string]string) ([]string, int) {
130 selected := make([]string, 0, len(siteIDs))
131 var skippedSelector int
132
133 for _, siteID := range siteIDs {
134 if !c.siteMatcher.MatchString(siteSelectorValue(siteID, siteNames[siteID])) {
135 skippedSelector++
136 continue
137 }
138 selected = append(selected, siteID)
139 }
140
141 return selected, skippedSelector
142 }
143
144 func siteSelectorValue(siteID, siteName string) string {
145 if name := strings.TrimSpace(siteName); name != "" {
146 return name
147 }
148 return strings.TrimSpace(siteID)
149 }
150
151 func (c *Collector) pruneUnselectedSites(sites map[string]*siteState, order *[]string) {
152 active := make(map[string]bool, len(c.discovery.siteIDs))
153 nextOrder := make([]string, 0, len(c.discovery.siteIDs))
154 for _, siteID := range c.discovery.siteIDs {
155 if sites[siteID] == nil {
156 continue
157 }
158 active[siteID] = true
159 nextOrder = append(nextOrder, siteID)
160 }
161 for siteID := range sites {
162 if !active[siteID] {
163 delete(sites, siteID)
164 }
165 }
166 *order = nextOrder
167 }
168
169 func (c *Collector) useCachedDiscoveryAfterRefreshFailure(now time.Time, force bool, err error) bool {
170 if force || len(c.discovery.siteIDs) == 0 {
171 return false
172 }
173
174 c.discovery.fetchedAt = now
175 errorClass := classifyCatoError(err)
176 c.warnRecoverable(warningKeyDiscoveryCache, errorClass, "entityLookup refresh failed; using cached discovery for %d site(s), error_class=%s", len(c.discovery.siteIDs), errorClass)
177 return true
178 }
179
180 func (c *Collector) collectSnapshot(ctx context.Context) (map[string]*siteState, []string, error) {
181 res, err := c.client.AccountSnapshot(ctx, c.AccountID, c.discovery.siteIDs)
182 if err != nil {
183 return nil, nil, err
184 }
185
186 sites, order := normalizeSnapshot(res, c.discovery.siteNames)
187 if len(order) == 0 {
188 return sites, order, nil
189 }
190
191 seen := make(map[string]bool, len(order))
192 for _, siteID := range order {
193 seen[siteID] = true
194 }
195 for _, siteID := range c.discovery.siteIDs {
196 if seen[siteID] {
197 continue
198 }
199 sites[siteID] = &siteState{
200 ID: siteID,
201 Name: siteDisplayName(siteID, c.discovery.siteNames, "", ""),
202 ConnectivityStatus: "unknown",
203 OperationalStatus: "unknown",
204 Interfaces: make(map[string]*interfaceState),
205 }
206 order = append(order, siteID)
207 }
208
209 return sites, order, nil
210 }
211
212 func seconds(v int) time.Duration {
213 return time.Duration(v) * time.Second
214 }