File size: 5,451 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 | package httpdriver
import (
"context"
"errors"
"fmt"
"net/http"
"github.com/samber/lo"
"github.com/openmeterio/openmeter/api"
"github.com/openmeterio/openmeter/openmeter/notification"
"github.com/openmeterio/openmeter/pkg/framework/commonhttp"
"github.com/openmeterio/openmeter/pkg/framework/transport/httptransport"
"github.com/openmeterio/openmeter/pkg/models"
"github.com/openmeterio/openmeter/pkg/pagination"
"github.com/openmeterio/openmeter/pkg/sortx"
)
type (
ListEventsRequest = notification.ListEventsInput
ListEventsResponse = api.NotificationEventPaginatedResponse
ListEventsParams = api.ListNotificationEventsParams
ListEventsHandler httptransport.HandlerWithArgs[ListEventsRequest, ListEventsResponse, ListEventsParams]
)
func (h *handler) ListEvents() ListEventsHandler {
return httptransport.NewHandlerWithArgs(
func(ctx context.Context, r *http.Request, params ListEventsParams) (ListEventsRequest, error) {
ns, err := h.resolveNamespace(ctx)
if err != nil {
return ListEventsRequest{}, fmt.Errorf("failed to resolve namespace: %w", err)
}
req := ListEventsRequest{
Namespaces: []string{ns},
Order: sortx.Order(lo.FromPtrOr(params.Order, api.SortOrderDESC)),
OrderBy: notification.OrderBy(lo.FromPtrOr(params.OrderBy, api.NotificationEventOrderByCreatedAt)),
Page: pagination.Page{
PageSize: lo.FromPtrOr(params.PageSize, notification.DefaultPageSize),
PageNumber: lo.FromPtrOr(params.Page, notification.DefaultPageNumber),
},
Subjects: lo.FromPtr(params.Subject),
Features: lo.FromPtr(params.Feature),
Rules: lo.FromPtr(params.Rule),
Channels: lo.FromPtr(params.Channel),
From: lo.FromPtr(params.From),
To: lo.FromPtr(params.To),
}
return req, nil
},
func(ctx context.Context, request ListEventsRequest) (ListEventsResponse, error) {
resp, err := h.service.ListEvents(ctx, request)
if err != nil {
return ListEventsResponse{}, fmt.Errorf("failed to list events: %w", err)
}
items := make([]api.NotificationEvent, 0, len(resp.Items))
for _, event := range resp.Items {
var item api.NotificationEvent
item, err = FromEvent(event)
if err != nil {
return ListEventsResponse{}, fmt.Errorf("failed to cast event: %w", err)
}
items = append(items, item)
}
return ListEventsResponse{
Items: items,
Page: resp.Page.PageNumber,
PageSize: resp.Page.PageSize,
TotalCount: resp.TotalCount,
}, nil
},
commonhttp.JSONResponseEncoderWithStatus[ListEventsResponse](http.StatusOK),
httptransport.AppendOptions(
h.options,
httptransport.WithOperationName("listNotificationEvents"),
httptransport.WithErrorEncoder(errorEncoder()),
)...,
)
}
type (
GetEventRequest = notification.GetEventInput
GetEventResponse = api.NotificationEvent
GetEventHandler httptransport.HandlerWithArgs[GetEventRequest, GetEventResponse, string]
)
func (h *handler) GetEvent() GetEventHandler {
return httptransport.NewHandlerWithArgs(
func(ctx context.Context, r *http.Request, eventID string) (GetEventRequest, error) {
ns, err := h.resolveNamespace(ctx)
if err != nil {
return GetEventRequest{}, fmt.Errorf("failed to resolve namespace: %w", err)
}
req := GetEventRequest{
Namespace: ns,
ID: eventID,
}
return req, nil
},
func(ctx context.Context, request GetEventRequest) (GetEventResponse, error) {
event, err := h.service.GetEvent(ctx, request)
if err != nil {
return GetEventResponse{}, fmt.Errorf("failed to get event: %w", err)
}
if event == nil {
return GetEventResponse{}, errors.New("failed to create test event: nil event returned")
}
return FromEvent(*event)
},
commonhttp.JSONResponseEncoderWithStatus[GetEventResponse](http.StatusOK),
httptransport.AppendOptions(
h.options,
httptransport.WithOperationName("getNotificationEvent"),
httptransport.WithErrorEncoder(errorEncoder()),
)...,
)
}
type (
ResendEventRequest = notification.ResendEventInput
ResendEventResponse = interface{}
ResendEventHandler httptransport.HandlerWithArgs[ResendEventRequest, ResendEventResponse, string]
)
func (h *handler) ResendEvent() ResendEventHandler {
return httptransport.NewHandlerWithArgs(
func(ctx context.Context, r *http.Request, eventID string) (ResendEventRequest, error) {
body := api.NotificationEventResendRequest{}
if err := commonhttp.JSONRequestBodyDecoder(r, &body); err != nil {
return ResendEventRequest{}, fmt.Errorf("field to decode resend event request: %w", err)
}
ns, err := h.resolveNamespace(ctx)
if err != nil {
return ResendEventRequest{}, fmt.Errorf("failed to resolve namespace: %w", err)
}
req := ResendEventRequest{
NamespacedID: models.NamespacedID{
Namespace: ns,
ID: eventID,
},
Channels: lo.FromPtr(body.Channels),
}
return req, nil
},
func(ctx context.Context, request ResendEventRequest) (ResendEventResponse, error) {
err := h.service.ResendEvent(ctx, request)
if err != nil {
return nil, fmt.Errorf("failed to resend event: %w", err)
}
return nil, nil
},
commonhttp.EmptyResponseEncoder[ResendEventResponse](http.StatusAccepted),
httptransport.AppendOptions(
h.options,
httptransport.WithOperationName("resendNotificationEvent"),
httptransport.WithErrorEncoder(errorEncoder()),
)...,
)
}
|