| package service |
|
|
| import ( |
| "context" |
| "errors" |
| "fmt" |
| "maps" |
|
|
| "github.com/samber/lo" |
|
|
| "github.com/openmeterio/openmeter/openmeter/customer" |
| "github.com/openmeterio/openmeter/openmeter/productcatalog" |
| "github.com/openmeterio/openmeter/openmeter/subscription" |
| subscriptionaddon "github.com/openmeterio/openmeter/openmeter/subscription/addon" |
| "github.com/openmeterio/openmeter/openmeter/subscription/patch" |
| subscriptionworkflow "github.com/openmeterio/openmeter/openmeter/subscription/workflow" |
| "github.com/openmeterio/openmeter/pkg/clock" |
| "github.com/openmeterio/openmeter/pkg/featuregate" |
| "github.com/openmeterio/openmeter/pkg/filter" |
| "github.com/openmeterio/openmeter/pkg/framework/transaction" |
| "github.com/openmeterio/openmeter/pkg/models" |
| "github.com/openmeterio/openmeter/pkg/pagination" |
| "github.com/openmeterio/openmeter/pkg/timeutil" |
| ) |
|
|
| func (s *service) CreateFromPlan(ctx context.Context, inp subscriptionworkflow.CreateSubscriptionWorkflowInput, plan subscription.Plan) (subscription.SubscriptionView, error) { |
| return transaction.Run(ctx, s.TransactionManager, func(ctx context.Context) (subscription.SubscriptionView, error) { |
| var def subscription.SubscriptionView |
|
|
| creditEnabled := featuregate.ContextResolver().Credits(ctx) |
| if !creditEnabled && plan.ToCreateSubscriptionPlanInput().SettlementMode == productcatalog.CreditOnlySettlementMode { |
| return def, models.NewGenericValidationError(errors.New("credits are not enabled on this deployment of OpenMeter")) |
| } |
|
|
| if err := s.lockCustomer(ctx, inp.CustomerID); err != nil { |
| return def, err |
| } |
|
|
| |
| cus, err := s.CustomerService.GetCustomer(ctx, customer.GetCustomerInput{ |
| CustomerID: &customer.CustomerID{ |
| Namespace: inp.Namespace, |
| ID: inp.CustomerID, |
| }, |
| }) |
| if err != nil { |
| return def, fmt.Errorf("failed to fetch customer: %w", err) |
| } |
|
|
| if cus != nil && cus.IsDeleted() { |
| return def, models.NewGenericPreConditionFailedError( |
| fmt.Errorf("customer is deleted [namespace=%s customer.id=%s]", cus.Namespace, cus.ID), |
| ) |
| } |
|
|
| if cus == nil { |
| return def, fmt.Errorf("unexpected nil customer") |
| } |
|
|
| if err := inp.Timing.ValidateForAction(subscription.SubscriptionActionCreate, nil); err != nil { |
| return def, fmt.Errorf("invalid timing: %w", err) |
| } |
|
|
| activeFrom, err := inp.Timing.Resolve() |
| if err != nil { |
| return def, fmt.Errorf("failed to resolve active from: %w", err) |
| } |
|
|
| |
| billingAnchor := lo.FromPtrOr(inp.BillingAnchor, activeFrom).UTC() |
|
|
| |
| spec, err := subscription.NewSpecFromPlan(plan, subscription.CreateSubscriptionCustomerInput{ |
| CustomerId: cus.ID, |
| Currency: plan.Currency(), |
| ActiveFrom: activeFrom, |
| MetadataModel: inp.MetadataModel, |
| Name: lo.CoalesceOrEmpty(inp.Name, plan.GetName()), |
| Description: inp.Description, |
| BillingAnchor: billingAnchor, |
| Annotations: inp.Annotations, |
| }) |
|
|
| if err := subscriptionworkflow.MapSubscriptionErrors(err); err != nil { |
| return def, fmt.Errorf("failed to create spec from plan: %w", err) |
| } |
|
|
| if err := spec.ValidateAlignment(); err != nil { |
| return def, err |
| } |
|
|
| |
| sub, err := s.Service.Create(ctx, inp.Namespace, spec) |
| if err != nil { |
| return def, fmt.Errorf("failed to create subscription: %w", err) |
| } |
|
|
| return s.Service.GetView(ctx, sub.NamespacedID) |
| }) |
| } |
|
|
| func (s *service) EditRunning(ctx context.Context, subscriptionID models.NamespacedID, customizations []subscription.Patch, timing subscription.Timing) (subscription.SubscriptionView, error) { |
| |
| return transaction.Run(ctx, s.TransactionManager, func(ctx context.Context) (subscription.SubscriptionView, error) { |
| |
| curr, err := s.Service.GetView(ctx, subscriptionID) |
| if err != nil { |
| return subscription.SubscriptionView{}, fmt.Errorf("failed to fetch subscription: %w", err) |
| } |
|
|
| adds, err := s.AddonService.List(ctx, subscriptionID.Namespace, subscriptionaddon.ListSubscriptionAddonsInput{ |
| SubscriptionID: subscriptionID.ID, |
| }) |
| if err != nil { |
| return subscription.SubscriptionView{}, fmt.Errorf("failed to list addons: %w", err) |
| } |
|
|
| if hasAddons(curr, adds.Items) { |
| return subscription.SubscriptionView{}, models.NewGenericForbiddenError(fmt.Errorf("subscription with addons cannot be edited")) |
| } |
|
|
| |
| |
| customizations = lo.Map(customizations, func(p subscription.Patch, _ int) subscription.Patch { |
| if ap, ok := p.(patch.PatchAddItem); ok { |
| if ap.CreateInput.CreateSubscriptionItemInput.Annotations == nil { |
| ap.CreateInput.CreateSubscriptionItemInput.Annotations = models.Annotations{} |
| } |
| _, _ = subscription.AnnotationParser.AddOwnerSubSystem(ap.CreateInput.CreateSubscriptionItemInput.Annotations, subscription.OwnerSubscriptionSubSystem) |
|
|
| subscriptionworkflow.AnnotationParser.SetUniquePatchID(ap.CreateInput.CreateSubscriptionItemInput.Annotations) |
|
|
| return ap |
| } |
|
|
| if ap, ok := p.(*patch.PatchAddItem); ok { |
| if ap.CreateInput.CreateSubscriptionItemInput.Annotations == nil { |
| ap.CreateInput.CreateSubscriptionItemInput.Annotations = models.Annotations{} |
| } |
| _, _ = subscription.AnnotationParser.AddOwnerSubSystem(ap.CreateInput.CreateSubscriptionItemInput.Annotations, subscription.OwnerSubscriptionSubSystem) |
|
|
| subscriptionworkflow.AnnotationParser.SetUniquePatchID(ap.CreateInput.CreateSubscriptionItemInput.Annotations) |
|
|
| return ap |
| } |
|
|
| return p |
| }) |
|
|
| |
| for i, p := range customizations { |
| if err := p.Validate(); err != nil { |
| return subscription.SubscriptionView{}, models.ErrorWithComponent(models.ComponentName(fmt.Sprintf("patch[%d]", i)), err) |
| } |
| } |
|
|
| |
| if err := timing.ValidateForAction(subscription.SubscriptionActionUpdate, &curr); err != nil { |
| return subscription.SubscriptionView{}, models.NewGenericValidationError(fmt.Errorf("invalid timing: %w", err)) |
| } |
|
|
| editTime, err := timing.ResolveForSpec(curr.Spec) |
| if err != nil { |
| return subscription.SubscriptionView{}, fmt.Errorf("failed to resolve timing: %w", err) |
| } |
|
|
| |
| spec := curr.AsSpec() |
|
|
| err = spec.ApplyMany(lo.Map(customizations, subscription.ToApplies), subscription.ApplyContext{ |
| CurrentTime: editTime, |
| }) |
| if err := subscriptionworkflow.MapSubscriptionErrors(err); err != nil { |
| return subscription.SubscriptionView{}, fmt.Errorf("failed to apply customizations: %w", err) |
| } |
|
|
| if err := spec.ValidateAlignment(); err != nil { |
| return subscription.SubscriptionView{}, err |
| } |
|
|
| sub, err := s.Service.Update(ctx, subscriptionID, spec) |
| if err != nil { |
| return subscription.SubscriptionView{}, fmt.Errorf("failed to update subscription: %w", err) |
| } |
|
|
| return s.Service.GetView(ctx, sub.NamespacedID) |
| }) |
| } |
|
|
| func (s *service) ChangeToPlan(ctx context.Context, subscriptionID models.NamespacedID, inp subscriptionworkflow.ChangeSubscriptionWorkflowInput, plan subscription.Plan) (subscription.Subscription, subscription.SubscriptionView, error) { |
| |
| type res struct { |
| curr subscription.Subscription |
| new subscription.SubscriptionView |
| } |
|
|
| |
| r, err := transaction.Run(ctx, s.TransactionManager, func(ctx context.Context) (res, error) { |
| creditEnabled := featuregate.ContextResolver().Credits(ctx) |
| if !creditEnabled && plan.ToCreateSubscriptionPlanInput().SettlementMode == productcatalog.CreditOnlySettlementMode { |
| return res{}, models.NewGenericValidationError(errors.New("credits are not enabled on this deployment of OpenMeter")) |
| } |
|
|
| |
| curr, err := s.Service.Cancel(ctx, subscriptionID, inp.Timing) |
| if err != nil { |
| return res{}, fmt.Errorf("failed to end current subscription: %w", err) |
| } |
|
|
| |
| verbatumTiming := subscription.Timing{ |
| Custom: curr.ActiveTo, |
| } |
|
|
| inp.Timing = verbatumTiming |
|
|
| |
| createInputAnnotations := models.Annotations{} |
| _, err = subscription.AnnotationParser.SetPreviousSubscriptionID(createInputAnnotations, curr.ID) |
| if err != nil { |
| return res{}, fmt.Errorf("failed to set previous subscription ID: %w", err) |
| } |
|
|
| |
| new, err := s.CreateFromPlan(ctx, subscriptionworkflow.CreateSubscriptionWorkflowInput{ |
| ChangeSubscriptionWorkflowInput: inp, |
| Namespace: curr.Namespace, |
| CustomerID: curr.CustomerId, |
| BillingAnchor: lo.ToPtr(lo.FromPtrOr(inp.BillingAnchor, curr.BillingAnchor)), |
| Annotations: createInputAnnotations, |
| }, plan) |
| if err != nil { |
| return res{}, fmt.Errorf("failed to create new subscription: %w", err) |
| } |
|
|
| |
| currAnnotations := curr.Annotations |
| if currAnnotations == nil { |
| currAnnotations = models.Annotations{} |
| } else { |
| currAnnotations = maps.Clone(currAnnotations) |
| } |
| currAnnotations, err = subscription.AnnotationParser.SetSupersedingSubscriptionID(currAnnotations, new.Subscription.ID) |
| if err != nil { |
| return res{}, fmt.Errorf("failed to set superseding subscription ID: %w", err) |
| } |
| updatedCurr, err := s.Service.UpdateAnnotations(ctx, curr.NamespacedID, currAnnotations) |
| if err != nil { |
| return res{}, fmt.Errorf("failed to update current subscription annotations: %w", err) |
| } |
| curr = *updatedCurr |
|
|
| |
| return res{curr, new}, nil |
| }) |
|
|
| return r.curr, r.new, err |
| } |
|
|
| func (s *service) Restore(ctx context.Context, subscriptionID models.NamespacedID) (subscription.Subscription, error) { |
| multiSubscriptionEnabled, err := s.FeatureFlags.IsFeatureEnabled(ctx, subscription.MultiSubscriptionEnabledFF) |
| if err != nil { |
| return subscription.Subscription{}, fmt.Errorf("failed to check if multi-subscription is enabled: %w", err) |
| } |
|
|
| if multiSubscriptionEnabled { |
| return subscription.Subscription{}, subscription.ErrRestoreSubscriptionNotAllowedForMultiSubscription |
| } |
|
|
| return transaction.Run(ctx, s.TransactionManager, func(ctx context.Context) (subscription.Subscription, error) { |
| now := clock.Now() |
|
|
| |
| sub, err := s.Service.GetView(ctx, subscriptionID) |
| if err != nil { |
| return subscription.Subscription{}, fmt.Errorf("failed to fetch subscription: %w", err) |
| } |
|
|
| |
| scheduled, err := pagination.CollectAll(ctx, pagination.NewPaginator(func(ctx context.Context, page pagination.Page) (pagination.Result[subscription.Subscription], error) { |
| return s.Service.List(ctx, subscription.ListSubscriptionsInput{ |
| CustomerID: &filter.FilterULID{FilterString: filter.FilterString{Eq: &sub.Subscription.CustomerId}}, |
| Namespaces: []string{sub.Subscription.Namespace}, |
| ActiveInPeriod: &timeutil.StartBoundedPeriod{From: now}, |
| Page: page, |
| }) |
| }), 1000) |
| if err != nil { |
| return subscription.Subscription{}, fmt.Errorf("failed to fetch scheduled subscriptions: %w", err) |
| } |
|
|
| |
| scheduled = lo.Filter(scheduled, func(s subscription.Subscription, _ int) bool { |
| return s.NamespacedID != subscriptionID |
| }) |
| |
| for _, sch := range scheduled { |
| if err := s.Service.Delete(ctx, sch.NamespacedID); err != nil { |
| return subscription.Subscription{}, fmt.Errorf("failed to delete scheduled subscription: %w", err) |
| } |
| } |
|
|
| |
| return s.Service.Continue(ctx, subscriptionID) |
| }) |
| } |
|
|