openmeter / customer /httpdriver /customer.go
Leon4gr45's picture
Upload folder using huggingface_hub (part 3)
1c4c66b verified
Raw
History Blame Contribute Delete
21.2 kB
package httpdriver
import (
"context"
"fmt"
"net/http"
"time"
"github.com/samber/lo"
"github.com/openmeterio/openmeter/api"
"github.com/openmeterio/openmeter/openmeter/customer"
"github.com/openmeterio/openmeter/openmeter/entitlement"
entitlementdriver "github.com/openmeterio/openmeter/openmeter/entitlement/driver"
"github.com/openmeterio/openmeter/openmeter/subscription"
"github.com/openmeterio/openmeter/pkg/clock"
"github.com/openmeterio/openmeter/pkg/defaultx"
"github.com/openmeterio/openmeter/pkg/filter"
"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"
)
// containsFilter wraps a *string query parameter in a FilterString with a
// Contains operator, preserving the v1 "case-insensitive partial match" semantics.
func containsFilter(value *string) *filter.FilterString {
if value == nil {
return nil
}
return &filter.FilterString{Contains: value}
}
// eqFilter wraps a *string query parameter in a FilterString with an Eq
// operator, preserving the v1 exact-match semantics.
func eqFilter(value *string) *filter.FilterString {
if value == nil {
return nil
}
return &filter.FilterString{Eq: value}
}
type (
ListCustomersResponse = pagination.Result[api.Customer]
ListCustomersParams = api.ListCustomersParams
ListCustomersRequest = customer.ListCustomersInput
ListCustomersHandler httptransport.HandlerWithArgs[ListCustomersRequest, ListCustomersResponse, ListCustomersParams]
)
// ListCustomers returns a handler for listing customers.
func (h *handler) ListCustomers() ListCustomersHandler {
return httptransport.NewHandlerWithArgs(
func(ctx context.Context, r *http.Request, params ListCustomersParams) (ListCustomersRequest, error) {
ns, err := h.resolveNamespace(ctx)
if err != nil {
return ListCustomersRequest{}, err
}
req := ListCustomersRequest{
Namespace: ns,
// Pagination
Page: pagination.Page{
PageSize: lo.FromPtrOr(params.PageSize, customer.DefaultPageSize),
PageNumber: lo.FromPtrOr(params.Page, customer.DefaultPageNumber),
},
// Order
OrderBy: string(defaultx.WithDefault(params.OrderBy, api.CustomerOrderByName)),
Order: sortx.Order(defaultx.WithDefault(params.Order, api.SortOrderASC)),
// Filters
Key: containsFilter(params.Key),
Name: containsFilter(params.Name),
PrimaryEmail: containsFilter(params.PrimaryEmail),
UsageAttributionSubjectKey: containsFilter(params.Subject),
PlanKey: eqFilter(params.PlanKey),
// Modifiers
IncludeDeleted: lo.FromPtrOr(params.IncludeDeleted, customer.IncludeDeleted),
// Expand
// TODO[v2]: disable expand of subscriptions by default, for now this is a breaking change
Expands: lo.Map(lo.FromPtrOr(params.Expand, api.QueryCustomerListExpand{api.CustomerExpandSubscriptions}), func(item api.CustomerExpand, _ int) customer.Expand {
return customer.Expand(item)
}),
}
if err := req.Page.Validate(); err != nil {
return ListCustomersRequest{}, err
}
return req, nil
},
func(ctx context.Context, request ListCustomersRequest) (ListCustomersResponse, error) {
resp, err := h.service.ListCustomers(ctx, request)
if err != nil {
return ListCustomersResponse{}, fmt.Errorf("failed to list customers: %w", err)
}
// Get the customer's subscriptions
var customerSubscriptions map[string][]subscription.Subscription
if len(resp.Items) > 0 {
customerIDs := lo.Map(resp.Items, func(item customer.Customer, _ int) string {
return item.ID
})
subscriptions, err := h.subscriptionService.List(ctx, subscription.ListSubscriptionsInput{
Namespaces: []string{request.Namespace},
CustomerID: &filter.FilterULID{FilterString: filter.FilterString{In: &customerIDs}},
ActiveAt: lo.ToPtr(time.Now()),
})
if err != nil {
return ListCustomersResponse{}, err
}
customerSubscriptions = lo.GroupBy(subscriptions.Items, func(item subscription.Subscription) string {
return item.CustomerId
})
}
// Map the customers to the API
return pagination.MapResultErr(resp, func(customer customer.Customer) (api.Customer, error) {
var item api.Customer
subs, ok := customerSubscriptions[customer.ID]
if !ok {
subs = []subscription.Subscription{}
}
item, err = CustomerToAPI(customer, subs, request.Expands)
if err != nil {
return item, fmt.Errorf("failed to cast customer customer: %w", err)
}
return item, nil
})
},
commonhttp.JSONResponseEncoderWithStatus[ListCustomersResponse](http.StatusOK),
httptransport.AppendOptions(
h.options,
httptransport.WithOperationName("listCustomers"),
)...,
)
}
type (
CreateCustomerRequest = customer.CreateCustomerInput
CreateCustomerResponse = api.Customer
CreateCustomerHandler httptransport.Handler[CreateCustomerRequest, CreateCustomerResponse]
)
// CreateCustomer returns a new httptransport.Handler for creating a customer.
func (h *handler) CreateCustomer() CreateCustomerHandler {
return httptransport.NewHandler(
func(ctx context.Context, r *http.Request) (CreateCustomerRequest, error) {
body := api.CustomerCreate{}
if err := commonhttp.JSONRequestBodyDecoder(r, &body); err != nil {
return CreateCustomerRequest{}, fmt.Errorf("field to decode create customer request: %w", err)
}
ns, err := h.resolveNamespace(ctx)
if err != nil {
return CreateCustomerRequest{}, err
}
req := CreateCustomerRequest{
Namespace: ns,
CustomerMutate: MapCustomerCreate(body),
}
return req, nil
},
func(ctx context.Context, request CreateCustomerRequest) (CreateCustomerResponse, error) {
customer, err := h.service.CreateCustomer(ctx, request)
if err != nil {
return CreateCustomerResponse{}, err
}
if customer == nil {
return CreateCustomerResponse{}, fmt.Errorf("failed to create customer")
}
return h.mapCustomerWithSubscriptionsToAPI(ctx, *customer, nil)
},
commonhttp.JSONResponseEncoderWithStatus[CreateCustomerResponse](http.StatusCreated),
httptransport.AppendOptions(
h.options,
httptransport.WithOperationName("createCustomer"),
)...,
)
}
type (
UpdateCustomerRequest struct {
Namespace string
CustomerIDOrKey string
CustomerMutate customer.CustomerMutate
}
UpdateCustomerResponse = api.Customer
UpdateCustomerHandler httptransport.HandlerWithArgs[UpdateCustomerRequest, UpdateCustomerResponse, string]
)
// UpdateCustomer returns a handler for updating a customer.
func (h *handler) UpdateCustomer() UpdateCustomerHandler {
return httptransport.NewHandlerWithArgs(
func(ctx context.Context, r *http.Request, customerIDOrKey string) (UpdateCustomerRequest, error) {
body := api.CustomerReplaceUpdate{}
if err := commonhttp.JSONRequestBodyDecoder(r, &body); err != nil {
return UpdateCustomerRequest{}, fmt.Errorf("field to decode update customer request: %w", err)
}
ns, err := h.resolveNamespace(ctx)
if err != nil {
return UpdateCustomerRequest{}, err
}
req := UpdateCustomerRequest{
Namespace: ns,
CustomerIDOrKey: customerIDOrKey,
CustomerMutate: MapCustomerReplaceUpdate(body),
}
return req, nil
},
func(ctx context.Context, request UpdateCustomerRequest) (UpdateCustomerResponse, error) {
// TODO: we should not allow key identifier for mutable operations
// Get the customer
cus, err := h.service.GetCustomer(ctx, customer.GetCustomerInput{
CustomerIDOrKey: &customer.CustomerIDOrKey{
IDOrKey: request.CustomerIDOrKey,
Namespace: request.Namespace,
},
})
if err != nil {
return UpdateCustomerResponse{}, err
}
if cus != nil && cus.IsDeleted() {
return UpdateCustomerResponse{},
models.NewGenericPreConditionFailedError(
fmt.Errorf("customer is deleted [namespace=%s customer.id=%s]", cus.Namespace, cus.ID),
)
}
customer, err := h.service.UpdateCustomer(ctx, customer.UpdateCustomerInput{
CustomerID: cus.GetID(),
CustomerMutate: request.CustomerMutate,
})
if err != nil {
return UpdateCustomerResponse{}, err
}
if customer == nil {
return UpdateCustomerResponse{}, fmt.Errorf("failed to update customer")
}
return h.mapCustomerWithSubscriptionsToAPI(ctx, *customer, nil)
},
commonhttp.JSONResponseEncoderWithStatus[UpdateCustomerResponse](http.StatusOK),
httptransport.AppendOptions(
h.options,
httptransport.WithOperationName("updateCustomer"),
)...,
)
}
type (
DeleteCustomerRequest struct {
Namespace string
CustomerIDOrKey string
}
DeleteCustomerResponse = interface{}
DeleteCustomerHandler httptransport.HandlerWithArgs[DeleteCustomerRequest, DeleteCustomerResponse, string]
)
// DeleteCustomer returns a handler for deleting a customer.
func (h *handler) DeleteCustomer() DeleteCustomerHandler {
return httptransport.NewHandlerWithArgs(
func(ctx context.Context, r *http.Request, customerIDOrKey string) (DeleteCustomerRequest, error) {
ns, err := h.resolveNamespace(ctx)
if err != nil {
return DeleteCustomerRequest{}, err
}
return DeleteCustomerRequest{
Namespace: ns,
CustomerIDOrKey: customerIDOrKey,
}, nil
},
func(ctx context.Context, request DeleteCustomerRequest) (DeleteCustomerResponse, error) {
// TODO: we should not allow key identifier for mutable operations
// Get the customer
cus, err := h.service.GetCustomer(ctx, customer.GetCustomerInput{
CustomerIDOrKey: &customer.CustomerIDOrKey{
IDOrKey: request.CustomerIDOrKey,
Namespace: request.Namespace,
},
})
if err != nil {
return DeleteCustomerRequest{}, err
}
if cus != nil && cus.IsDeleted() {
return DeleteCustomerRequest{},
models.NewGenericPreConditionFailedError(
fmt.Errorf("customer is deleted [namespace=%s customer.id=%s]", cus.Namespace, cus.ID),
)
}
err = h.service.DeleteCustomer(ctx, cus.GetID())
if err != nil {
return nil, err
}
return nil, nil
},
commonhttp.EmptyResponseEncoder[DeleteCustomerResponse](http.StatusNoContent),
httptransport.AppendOptions(
h.options,
httptransport.WithOperationName("deleteCustomer"),
)...,
)
}
type (
GetCustomerRequest = customer.GetCustomerInput
GetCustomerResponse = api.Customer
GetCustomerHandler httptransport.HandlerWithArgs[GetCustomerRequest, GetCustomerResponse, GetCustomerParams]
)
type GetCustomerParams struct {
CustomerIDOrKey string
api.GetCustomerParams
}
// GetCustomer returns a handler for getting a customer.
func (h *handler) GetCustomer() GetCustomerHandler {
return httptransport.NewHandlerWithArgs(
func(ctx context.Context, r *http.Request, params GetCustomerParams) (GetCustomerRequest, error) {
ns, err := h.resolveNamespace(ctx)
if err != nil {
return GetCustomerRequest{}, err
}
return GetCustomerRequest{
CustomerIDOrKey: &customer.CustomerIDOrKey{
Namespace: ns,
IDOrKey: params.CustomerIDOrKey,
},
Expands: lo.Map(lo.FromPtrOr(params.Expand, api.QueryCustomerListExpand{api.CustomerExpandSubscriptions}), func(item api.CustomerExpand, _ int) customer.Expand {
return customer.Expand(item)
}),
}, nil
},
func(ctx context.Context, request GetCustomerRequest) (GetCustomerResponse, error) {
// Get the customer
cus, err := h.service.GetCustomer(ctx, request)
if err != nil {
return GetCustomerResponse{}, err
}
if cus == nil {
return GetCustomerResponse{}, fmt.Errorf("failed to get customer")
}
return h.mapCustomerWithSubscriptionsToAPI(ctx, *cus, request.Expands)
},
commonhttp.JSONResponseEncoderWithStatus[GetCustomerResponse](http.StatusOK),
httptransport.AppendOptions(
h.options,
httptransport.WithOperationName("getCustomer"),
)...,
)
}
type (
GetCustomerEntitlementValueRequest struct {
customer.GetEntitlementValueInput
At time.Time
}
GetCustomerEntitlementValueResponse = api.EntitlementValue
GetCustomerEntitlementValueParams = struct {
CustomerIDOrKey string
FeatureKey string
Time *time.Time
}
GetCustomerEntitlementValueHandler httptransport.HandlerWithArgs[GetCustomerEntitlementValueRequest, GetCustomerEntitlementValueResponse, GetCustomerEntitlementValueParams]
)
type (
GetCustomerEntitlementValueV2Response = EntitlementValueV2
GetCustomerEntitlementValueV2Handler httptransport.HandlerWithArgs[GetCustomerEntitlementValueRequest, GetCustomerEntitlementValueV2Response, GetCustomerEntitlementValueParams]
)
// GetCustomerEntitlementValue returns a handler for getting a customer.
func (h *handler) GetCustomerEntitlementValue() GetCustomerEntitlementValueHandler {
return httptransport.NewHandlerWithArgs(
func(ctx context.Context, r *http.Request, params GetCustomerEntitlementValueParams) (GetCustomerEntitlementValueRequest, error) {
ns, err := h.resolveNamespace(ctx)
if err != nil {
return GetCustomerEntitlementValueRequest{}, err
}
// Get the customer
cus, err := h.service.GetCustomer(ctx, customer.GetCustomerInput{
CustomerIDOrKey: &customer.CustomerIDOrKey{
IDOrKey: params.CustomerIDOrKey,
Namespace: ns,
},
})
if err != nil {
return GetCustomerEntitlementValueRequest{}, err
}
if cus != nil && cus.IsDeleted() {
return GetCustomerEntitlementValueRequest{},
models.NewGenericPreConditionFailedError(
fmt.Errorf("customer is deleted [namespace=%s customer.id=%s]", cus.Namespace, cus.ID),
)
}
return GetCustomerEntitlementValueRequest{
GetEntitlementValueInput: customer.GetEntitlementValueInput{
FeatureKey: params.FeatureKey,
CustomerID: cus.GetID(),
},
At: defaultx.WithDefault(params.Time, clock.Now()),
}, nil
},
func(ctx context.Context, request GetCustomerEntitlementValueRequest) (GetCustomerEntitlementValueResponse, error) {
val, err := h.entitlementService.GetEntitlementValue(ctx, request.CustomerID.Namespace, request.CustomerID.ID, request.FeatureKey, request.At)
if err != nil {
if _, ok := lo.ErrorsAs[*entitlement.NotFoundError](err); ok {
val = &entitlement.NoAccessValue{}
err = nil
}
}
if err != nil {
return GetCustomerEntitlementValueResponse{}, err
}
return entitlementdriver.MapEntitlementValueToAPI(val)
},
commonhttp.JSONResponseEncoderWithStatus[GetCustomerEntitlementValueResponse](http.StatusOK),
httptransport.AppendOptions(
h.options,
httptransport.WithOperationName("getCustomer"),
)...,
)
}
func (h *handler) GetCustomerEntitlementValueV2() GetCustomerEntitlementValueV2Handler {
return httptransport.NewHandlerWithArgs(
func(ctx context.Context, r *http.Request, params GetCustomerEntitlementValueParams) (GetCustomerEntitlementValueRequest, error) {
ns, err := h.resolveNamespace(ctx)
if err != nil {
return GetCustomerEntitlementValueRequest{}, err
}
cus, err := h.service.GetCustomer(ctx, customer.GetCustomerInput{
CustomerIDOrKey: &customer.CustomerIDOrKey{
IDOrKey: params.CustomerIDOrKey,
Namespace: ns,
},
})
if err != nil {
return GetCustomerEntitlementValueRequest{}, err
}
if cus != nil && cus.IsDeleted() {
return GetCustomerEntitlementValueRequest{},
models.NewGenericPreConditionFailedError(
fmt.Errorf("customer is deleted [namespace=%s customer.id=%s]", cus.Namespace, cus.ID),
)
}
return GetCustomerEntitlementValueRequest{
GetEntitlementValueInput: customer.GetEntitlementValueInput{
FeatureKey: params.FeatureKey,
CustomerID: cus.GetID(),
},
At: defaultx.WithDefault(params.Time, clock.Now()),
}, nil
},
func(ctx context.Context, request GetCustomerEntitlementValueRequest) (GetCustomerEntitlementValueV2Response, error) {
val, err := h.entitlementService.GetEntitlementValue(ctx, request.CustomerID.Namespace, request.CustomerID.ID, request.FeatureKey, request.At)
if err != nil {
if _, ok := lo.ErrorsAs[*entitlement.NotFoundError](err); ok {
val = &entitlement.NoAccessValue{}
err = nil
}
}
if err != nil {
return GetCustomerEntitlementValueV2Response{}, err
}
return MapEntitlementValueToAPIV2(val)
},
commonhttp.JSONResponseEncoderWithStatus[GetCustomerEntitlementValueV2Response](http.StatusOK),
httptransport.AppendOptions(
h.options,
httptransport.WithOperationName("getCustomerEntitlementValueV2"),
)...,
)
}
type (
GetCustomerAccessRequest = customer.GetCustomerInput
GetCustomerAccessResponse = api.CustomerAccess
GetCustomerAccessParams = struct {
CustomerIDOrKey string
}
GetCustomerAccessHandler httptransport.HandlerWithArgs[GetCustomerAccessRequest, GetCustomerAccessResponse, GetCustomerAccessParams]
)
type (
GetCustomerAccessV2Response = CustomerAccessV2
GetCustomerAccessV2Handler httptransport.HandlerWithArgs[GetCustomerAccessRequest, GetCustomerAccessV2Response, GetCustomerAccessParams]
)
// GetCustomerAccess returns a handler for getting a customer access.
func (h *handler) GetCustomerAccess() GetCustomerAccessHandler {
return httptransport.NewHandlerWithArgs(
func(ctx context.Context, r *http.Request, params GetCustomerAccessParams) (GetCustomerAccessRequest, error) {
ns, err := h.resolveNamespace(ctx)
if err != nil {
return GetCustomerAccessRequest{}, err
}
return GetCustomerAccessRequest{
CustomerIDOrKey: &customer.CustomerIDOrKey{
Namespace: ns,
IDOrKey: params.CustomerIDOrKey,
},
}, nil
},
func(ctx context.Context, request GetCustomerAccessRequest) (GetCustomerAccessResponse, error) {
cus, err := h.service.GetCustomer(ctx, customer.GetCustomerInput{
CustomerIDOrKey: request.CustomerIDOrKey,
})
if err != nil {
return GetCustomerAccessResponse{}, err
}
if cus != nil && cus.IsDeleted() {
return GetCustomerAccessResponse{},
models.NewGenericPreConditionFailedError(
fmt.Errorf("customer is deleted [namespace=%s customer.id=%s]", cus.Namespace, cus.ID),
)
}
access, err := h.entitlementService.GetAccess(ctx, cus.Namespace, cus.ID)
if err != nil {
return GetCustomerAccessResponse{}, err
}
apiAccess, err := MapAccessToAPI(access)
if err != nil {
return GetCustomerAccessResponse{}, err
}
return apiAccess, nil
},
commonhttp.JSONResponseEncoderWithStatus[GetCustomerAccessResponse](http.StatusOK),
httptransport.AppendOptions(
h.options,
httptransport.WithOperationName("getCustomerAccess"),
)...,
)
}
func (h *handler) GetCustomerAccessV2() GetCustomerAccessV2Handler {
return httptransport.NewHandlerWithArgs(
func(ctx context.Context, r *http.Request, params GetCustomerAccessParams) (GetCustomerAccessRequest, error) {
ns, err := h.resolveNamespace(ctx)
if err != nil {
return GetCustomerAccessRequest{}, err
}
return GetCustomerAccessRequest{
CustomerIDOrKey: &customer.CustomerIDOrKey{
Namespace: ns,
IDOrKey: params.CustomerIDOrKey,
},
}, nil
},
func(ctx context.Context, request GetCustomerAccessRequest) (GetCustomerAccessV2Response, error) {
cus, err := h.service.GetCustomer(ctx, customer.GetCustomerInput{
CustomerIDOrKey: request.CustomerIDOrKey,
})
if err != nil {
return GetCustomerAccessV2Response{}, err
}
if cus != nil && cus.IsDeleted() {
return GetCustomerAccessV2Response{},
models.NewGenericPreConditionFailedError(
fmt.Errorf("customer is deleted [namespace=%s customer.id=%s]", cus.Namespace, cus.ID),
)
}
access, err := h.entitlementService.GetAccess(ctx, cus.Namespace, cus.ID)
if err != nil {
return GetCustomerAccessV2Response{}, err
}
apiAccess, err := MapAccessToAPIV2(access)
if err != nil {
return GetCustomerAccessV2Response{}, err
}
return apiAccess, nil
},
commonhttp.JSONResponseEncoderWithStatus[GetCustomerAccessV2Response](http.StatusOK),
httptransport.AppendOptions(
h.options,
httptransport.WithOperationName("getCustomerAccessV2"),
)...,
)
}
// mapCustomerWithSubscriptionsToAPI maps a customer to the API with its subscriptions.
func (h *handler) mapCustomerWithSubscriptionsToAPI(ctx context.Context, cust customer.Customer, expand []customer.Expand) (api.Customer, error) {
if !lo.Contains(expand, customer.ExpandSubscriptions) {
return CustomerToAPI(cust, []subscription.Subscription{}, expand)
}
// Get the customer's subscriptions
subscriptions, err := h.subscriptionService.List(ctx, subscription.ListSubscriptionsInput{
Namespaces: []string{cust.Namespace},
CustomerID: &filter.FilterULID{FilterString: filter.FilterString{Eq: &cust.ID}},
ActiveAt: lo.ToPtr(time.Now()),
})
if err != nil {
return GetCustomerResponse{}, err
}
// Map the customer to the API
return CustomerToAPI(cust, subscriptions.Items, expand)
}