-
Notifications
You must be signed in to change notification settings - Fork 53
/
aggregation.go
171 lines (158 loc) · 4.48 KB
/
aggregation.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
package cortex
import (
"bytes"
"compress/flate"
"compress/gzip"
"encoding/json"
"fmt"
"io"
"net/http"
"sync"
"github.com/andybalholm/brotli"
"github.com/gin-gonic/gin"
managementv1 "github.com/rancher/opni/pkg/apis/management/v1"
"github.com/rancher/opni/pkg/logger"
"github.com/rancher/opni/pkg/rbac"
"go.uber.org/zap"
)
type DataFormat string
const (
// Rule data formatted as a single YAML document containing yaml-encoded
// []rulefmt.RuleGroup keyed by tenant ID:
// <tenantID>:
// - name: ...
// rules: [...]
NamespaceKeyedYAML DataFormat = "namespace-keyed-yaml"
// Rule data formatted as JSON containing the prometheus response metadata
// and a list of rule groups, each of which has a field "file" containing
// the tenant ID:
// {"status":"success","data":{"groups":["file":"<tenantID>", ...]}}
PrometheusRuleGroupsJSON DataFormat = "prometheus-rule-groups-json"
)
type MultiTenantRuleAggregator struct {
client managementv1.ManagementClient
cortexClient *http.Client
headerCodec rbac.HeaderCodec
bufferPool *sync.Pool
logger *zap.SugaredLogger
format DataFormat
}
func NewMultiTenantRuleAggregator(
client managementv1.ManagementClient,
cortexClient *http.Client,
headerCodec rbac.HeaderCodec,
format DataFormat,
) *MultiTenantRuleAggregator {
pool := &sync.Pool{
New: func() any {
return bytes.NewBuffer(make([]byte, 0, 1024*1024)) // 1MB
},
}
return &MultiTenantRuleAggregator{
client: client,
cortexClient: cortexClient,
headerCodec: headerCodec,
bufferPool: pool,
logger: logger.New().Named("aggregation"),
format: format,
}
}
type rawPromJSONData struct {
Status string `json:"status"`
Data struct {
Groups []json.RawMessage `json:"groups"`
} `json:"data"`
ErrorType string `json:"errorType"`
Error string `json:"error"`
}
func (a *MultiTenantRuleAggregator) Handle(c *gin.Context) {
ids := rbac.AuthorizedClusterIDs(c)
a.logger.With(
"request", c.FullPath(),
).Debugf("aggregating query over %d tenants", len(ids))
buf := a.bufferPool.Get().(*bytes.Buffer)
switch a.format {
case NamespaceKeyedYAML:
buf.WriteString("---\n")
case PrometheusRuleGroupsJSON:
buf.WriteString(`{"status":"success","data":{"groups":[`)
}
// Repeat the request for each cluster ID, sending one ID at a time, and
// return the aggregated results.
groupCount := 0
for _, id := range ids {
req := c.Copy().Request
req.Header.Set(a.headerCodec.Key(), a.headerCodec.Encode([]string{id}))
resp, err := a.cortexClient.Do(req)
if err != nil {
a.logger.With(
"request", c.FullPath(),
"error", err,
).Error("error querying cortex")
c.AbortWithError(http.StatusInternalServerError, err)
return
}
defer resp.Body.Close()
if code := resp.StatusCode; code != http.StatusOK {
if code == http.StatusNotFound {
// cortex will report 404 if there are no groups
continue
}
return
}
var bodyReader io.ReadCloser
switch encoding := req.Header.Get("Content-Encoding"); encoding {
case "gzip":
bodyReader, _ = gzip.NewReader(resp.Body)
case "brotli":
bodyReader = io.NopCloser(brotli.NewReader(resp.Body))
case "deflate":
bodyReader = flate.NewReader(resp.Body)
case "":
bodyReader = resp.Body
default:
c.AbortWithError(http.StatusBadRequest, fmt.Errorf("unsupported Content-Encoding: %s", encoding))
return
}
body, err := io.ReadAll(bodyReader)
if err != nil {
c.AbortWithError(http.StatusInternalServerError, err)
return
}
switch a.format {
case NamespaceKeyedYAML:
// we can simply concatenate the responses
buf.Write(body)
case PrometheusRuleGroupsJSON:
rawData := rawPromJSONData{}
if err := json.Unmarshal(body, &rawData); err != nil {
c.AbortWithError(http.StatusInternalServerError, fmt.Errorf("error parsing cortex response: %s", err))
return
}
if rawData.Status != "success" {
a.logger.With(
"request", c.FullPath(),
"status", rawData.Status,
"errorType", rawData.ErrorType,
"error", rawData.Error,
).Error("error response from prometheus")
continue
}
for _, group := range rawData.Data.Groups {
raw, _ := group.MarshalJSON()
if groupCount > 0 {
buf.WriteString(",")
}
groupCount++
buf.Write(raw)
}
}
}
switch a.format {
case NamespaceKeyedYAML:
c.Data(http.StatusOK, "application/yaml", buf.Bytes())
case PrometheusRuleGroupsJSON:
buf.WriteString(`]},"errorType":"","error":""}`)
c.Data(http.StatusOK, "application/json", buf.Bytes())
}
}