| 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" |
| ) |
|
|
| |
| 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] |
| ) |
|
|
| |
| 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) { |
| |
| 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) |
| } |
|
|
| |
| params, err := h.toQueryParamsFromRequest( |
| ctx, |
| meter, |
| |
| 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) |
| } |
|
|
| |
| subjectKeys := getSubjectsFromQueryResult(rows) |
|
|
| var subjectsByKey map[string]subject.Subject |
|
|
| |
| 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"), |
| )..., |
| ) |
| } |
|
|
| |
| 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) { |
| |
| 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) |
| } |
|
|
| |
| 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) |
| } |
|
|
| |
| subjectKeys := getSubjectsFromQueryResult(rows) |
|
|
| var subjectsByKey map[string]subject.Subject |
|
|
| |
| 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"), |
| )..., |
| ) |
| } |
|
|
| |
| func getSubjectsFromQueryResult(rows []meter.MeterQueryRow) []string { |
| |
| subjects := []string{} |
| for _, row := range rows { |
| if row.Subject == nil { |
| continue |
| } |
|
|
| subjects = append(subjects, *row.Subject) |
| } |
|
|
| |
| 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 |
| } |
|
|
| |
| func (a *queryMeterCSVResult) Records() [][]string { |
| records := [][]string{} |
|
|
| hasSubjectHeader := false |
|
|
| |
| groupByKeys := []string{} |
| for _, k := range a.queryGroupBy { |
| if k == "subject" { |
| hasSubjectHeader = true |
| continue |
| } |
| groupByKeys = append(groupByKeys, k) |
| } |
|
|
| |
| headers := []string{"window_start", "window_end"} |
|
|
| |
| if hasSubjectHeader { |
| headers = append(headers, "subject") |
| } |
|
|
| |
| enhanceSubject := hasSubjectHeader && len(a.subjectsByKey) > 0 |
|
|
| if enhanceSubject { |
| headers = append(headers, "subject_display_name") |
| } |
|
|
| |
| if len(groupByKeys) > 0 { |
| headers = append(headers, groupByKeys...) |
| } |
| headers = append(headers, "value") |
|
|
| records = append(records, headers) |
|
|
| |
| for _, row := range a.rows { |
| data := []string{row.WindowStart.Format(time.RFC3339), row.WindowEnd.Format(time.RFC3339)} |
|
|
| |
| if hasSubjectHeader { |
| data = append(data, lo.FromPtrOr(row.Subject, "")) |
|
|
| |
| 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 |
| } |
|
|
| |
| func (a *queryMeterCSVResult) FileName() string { |
| return a.meterSlug |
| } |
|
|