File size: 8,198 Bytes
d6f631f | 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 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 | package httpdriver
import (
"context"
"fmt"
"net/http"
"time"
"github.com/samber/lo"
"github.com/openmeterio/openmeter/api"
"github.com/openmeterio/openmeter/openmeter/meter"
"github.com/openmeterio/openmeter/openmeter/subject"
"github.com/openmeterio/openmeter/pkg/framework/commonhttp"
"github.com/openmeterio/openmeter/pkg/framework/transport/httptransport"
)
// QueryMeterCSVResult is a CSV response for query meter.
var _ QueryMeterCSVResponse = (*queryMeterCSVResult)(nil)
type (
QueryMeterCSVParams = QueryMeterParams
QueryMeterCSVRequest = QueryMeterRequest
QueryMeterCSVResponse = commonhttp.CSVResponse
QueryMeterCSVHandler httptransport.HandlerWithArgs[QueryMeterCSVRequest, QueryMeterCSVResponse, QueryMeterCSVParams]
)
type (
QueryMeterPostCSVParams = QueryMeterPostParams
QueryMeterPostCSVRequest = QueryMeterPostRequest
QueryMeterPostCSVResponse = commonhttp.CSVResponse
QueryMeterPostCSVHandler httptransport.HandlerWithArgs[QueryMeterPostCSVRequest, QueryMeterPostCSVResponse, QueryMeterPostCSVParams]
)
// QueryMeterCSV returns a handler for query meter.
func (h *handler) QueryMeterCSV() QueryMeterCSVHandler {
return httptransport.NewHandlerWithArgs(
func(ctx context.Context, r *http.Request, params QueryMeterCSVParams) (QueryMeterCSVRequest, error) {
ns, err := h.resolveNamespace(ctx)
if err != nil {
return QueryMeterCSVRequest{}, err
}
return QueryMeterCSVRequest{
namespace: ns,
idOrSlug: params.IdOrSlug,
params: params.QueryMeterParams,
}, nil
},
func(ctx context.Context, request QueryMeterCSVRequest) (QueryMeterCSVResponse, error) {
// Get meter
meter, err := h.meterService.GetMeterByIDOrSlug(ctx, meter.GetMeterInput{
Namespace: request.namespace,
IDOrSlug: request.idOrSlug,
})
if err != nil {
return nil, fmt.Errorf("failed to get meter: %w", err)
}
// Query meter
params, err := h.toQueryParamsFromRequest(
ctx,
meter,
// Convert the POST request body to a GET request params
ToRequestFromQueryParamsPOSTBody(request.params),
)
if err != nil {
return nil, fmt.Errorf("failed to construct query meter params: %w", err)
}
rows, err := h.streaming.QueryMeter(ctx, request.namespace, meter, params)
if err != nil {
return nil, fmt.Errorf("failed to query meter: %w", err)
}
// Collect subjects from query results if any
subjectKeys := getSubjectsFromQueryResult(rows)
var subjectsByKey map[string]subject.Subject
// If there are subjects get the display names
if len(subjectKeys) > 0 {
subjects, err := h.subjectService.List(ctx, request.namespace, subject.ListParams{
Keys: subjectKeys,
})
if err != nil {
return nil, fmt.Errorf("failed to get subjects: %w", err)
}
subjectsByKey = lo.KeyBy(subjects.Items, func(s subject.Subject) string {
return s.Key
})
}
response := NewQueryMeterCSVResult(meter.Key, params.GroupBy, rows, subjectsByKey)
return response, nil
},
commonhttp.CSVResponseEncoder[QueryMeterCSVResponse],
httptransport.AppendOptions(
h.options,
httptransport.WithOperationName("queryMeterCSV"),
)...,
)
}
// QueryMeterPostCSV returns a handler for query meter via POST with CSV response.
func (h *handler) QueryMeterPostCSV() QueryMeterPostCSVHandler {
return httptransport.NewHandlerWithArgs(
func(ctx context.Context, r *http.Request, meterIdOrSlug QueryMeterPostCSVParams) (QueryMeterPostCSVRequest, error) {
ns, err := h.resolveNamespace(ctx)
if err != nil {
return QueryMeterPostCSVRequest{}, err
}
var request api.QueryMeterPostJSONRequestBody
if err := commonhttp.JSONRequestBodyDecoder(r, &request); err != nil {
return QueryMeterPostCSVRequest{}, fmt.Errorf("failed to decode request body: %w", err)
}
return QueryMeterPostCSVRequest{
namespace: ns,
idOrSlug: meterIdOrSlug,
params: request,
}, nil
},
func(ctx context.Context, request QueryMeterPostCSVRequest) (QueryMeterPostCSVResponse, error) {
// Get meter
meter, err := h.meterService.GetMeterByIDOrSlug(ctx, meter.GetMeterInput{
Namespace: request.namespace,
IDOrSlug: request.idOrSlug,
})
if err != nil {
return nil, fmt.Errorf("failed to get meter: %w", err)
}
// Query meter
params, err := h.toQueryParamsFromRequest(ctx, meter, request.params)
if err != nil {
return nil, fmt.Errorf("failed to construct query meter params: %w", err)
}
rows, err := h.streaming.QueryMeter(ctx, request.namespace, meter, params)
if err != nil {
return nil, fmt.Errorf("failed to query meter: %w", err)
}
// Collect subjects from query results if any
subjectKeys := getSubjectsFromQueryResult(rows)
var subjectsByKey map[string]subject.Subject
// If there are subjects get the display names
if len(subjectKeys) > 0 {
subjects, err := h.subjectService.List(ctx, request.namespace, subject.ListParams{
Keys: subjectKeys,
})
if err != nil {
return nil, fmt.Errorf("failed to get subjects: %w", err)
}
subjectsByKey = lo.KeyBy(subjects.Items, func(s subject.Subject) string {
return s.Key
})
}
response := NewQueryMeterCSVResult(meter.Key, params.GroupBy, rows, subjectsByKey)
return response, nil
},
commonhttp.CSVResponseEncoder[QueryMeterPostCSVResponse],
httptransport.AppendOptions(
h.options,
httptransport.WithOperationName("queryMeterPostCSV"),
)...,
)
}
// getSubjectsFromQueryResult returns the subjects from a query result.
func getSubjectsFromQueryResult(rows []meter.MeterQueryRow) []string {
// Collect subjects from query results if any
subjects := []string{}
for _, row := range rows {
if row.Subject == nil {
continue
}
subjects = append(subjects, *row.Subject)
}
// Deduplicate subjects
subjects = lo.Uniq(subjects)
return subjects
}
func NewQueryMeterCSVResult(meterSlug string, queryGroupBy []string, rows []meter.MeterQueryRow, subjectsByKey map[string]subject.Subject) QueryMeterCSVResponse {
return &queryMeterCSVResult{
meterSlug: meterSlug,
queryGroupBy: queryGroupBy,
rows: rows,
subjectsByKey: subjectsByKey,
}
}
type queryMeterCSVResult struct {
meterSlug string
queryGroupBy []string
rows []meter.MeterQueryRow
subjectsByKey map[string]subject.Subject
}
// Records returns the CSV records.
func (a *queryMeterCSVResult) Records() [][]string {
records := [][]string{}
hasSubjectHeader := false
// Filter out the subject from the group by keys
groupByKeys := []string{}
for _, k := range a.queryGroupBy {
if k == "subject" {
hasSubjectHeader = true
continue
}
groupByKeys = append(groupByKeys, k)
}
// CSV headers
headers := []string{"window_start", "window_end"}
// Add subject header if present
if hasSubjectHeader {
headers = append(headers, "subject")
}
// Add subject display name header if present
enhanceSubject := hasSubjectHeader && len(a.subjectsByKey) > 0
if enhanceSubject {
headers = append(headers, "subject_display_name")
}
// Add group by headers if present
if len(groupByKeys) > 0 {
headers = append(headers, groupByKeys...)
}
headers = append(headers, "value")
records = append(records, headers)
// CSV data
for _, row := range a.rows {
data := []string{row.WindowStart.Format(time.RFC3339), row.WindowEnd.Format(time.RFC3339)}
// Add subject if header is present
if hasSubjectHeader {
data = append(data, lo.FromPtrOr(row.Subject, ""))
// Add display name if available
if enhanceSubject && row.Subject != nil {
subject, ok := a.subjectsByKey[*row.Subject]
if ok {
data = append(data, lo.FromPtrOr(subject.DisplayName, ""))
} else {
data = append(data, "")
}
}
}
for _, k := range groupByKeys {
var groupByValue string
if row.GroupBy[k] != nil {
groupByValue = *row.GroupBy[k]
}
data = append(data, groupByValue)
}
data = append(data, fmt.Sprintf("%f", row.Value))
records = append(records, data)
}
return records
}
// FileName returns the CSV file name.
func (a *queryMeterCSVResult) FileName() string {
return a.meterSlug
}
|