openmeter / app /common /notification.go
Leon4gr45's picture
Upload folder using huggingface_hub (part 4)
1f10f31 verified
Raw
History Blame Contribute Delete
4.79 kB
package common
import (
"fmt"
"log/slog"
"math/rand/v2"
"github.com/google/wire"
svix "github.com/svix/svix-webhooks/go"
"go.opentelemetry.io/otel/trace"
"github.com/openmeterio/openmeter/app/config"
entdb "github.com/openmeterio/openmeter/openmeter/ent/db"
"github.com/openmeterio/openmeter/openmeter/notification"
notificationadapter "github.com/openmeterio/openmeter/openmeter/notification/adapter"
"github.com/openmeterio/openmeter/openmeter/notification/eventhandler"
notificationservice "github.com/openmeterio/openmeter/openmeter/notification/service"
notificationwebhook "github.com/openmeterio/openmeter/openmeter/notification/webhook"
webhooknoop "github.com/openmeterio/openmeter/openmeter/notification/webhook/noop"
webhooksvix "github.com/openmeterio/openmeter/openmeter/notification/webhook/svix"
"github.com/openmeterio/openmeter/openmeter/productcatalog/feature"
"github.com/openmeterio/openmeter/pkg/framework/pgdriver"
"github.com/openmeterio/openmeter/pkg/pglockx"
)
var Notification = wire.NewSet(
NewNotificationAdapter,
NewNotificationService,
NewNotificationWebhookHandler,
NewNotificationEventHandler,
)
// NotificationService is a wire set for the notification service, it can be used at
// places where only the service is required without svix and event handler.
var NotificationService = wire.NewSet(
NewNotificationAdapter,
NewNotificationService,
NewNoopNotificationWebhookHandler,
)
func NewNotificationAdapter(
logger *slog.Logger,
db *entdb.Client,
) (notification.Repository, error) {
adapter, err := notificationadapter.New(notificationadapter.Config{
Client: db,
Logger: logger,
})
if err != nil {
return nil, fmt.Errorf("failed to initialize notification adapter: %w", err)
}
return adapter, nil
}
func NewNotificationEventHandler(
config config.NotificationConfiguration,
logger *slog.Logger,
tracer trace.Tracer,
adapter notification.Repository,
webhook notificationwebhook.Handler,
driver *pgdriver.Driver,
) (notification.EventHandler, error) {
config.Lock.Owner = fmt.Sprintf("notification.event_handler-%v", rand.Int())
config.Lock.HeartbeatInterval = max(pglockx.DefaultHeartbeatInterval, config.Lock.HeartbeatInterval)
config.Lock.LeaseTime = max(config.Lock.HeartbeatInterval*2, config.Lock.LeaseTime)
logger.Debug("initializing notification lock client",
"lock.leaseTime", config.Lock.LeaseTime, "lock.heartbeatInterval", config.Lock.HeartbeatInterval, "lock.owner", config.Lock.Owner)
lockClient, err := pglockx.New(driver.DB(), config.Lock)
if err != nil {
return nil, fmt.Errorf("failed to initialize notification lock client: %w", err)
}
eventHandler, err := eventhandler.New(eventhandler.Config{
Repository: adapter,
Webhook: webhook,
Logger: logger,
Tracer: tracer,
ReconcileInterval: config.ReconcileInterval,
SendingTimeout: config.SendingTimeout,
PendingTimeout: config.PendingTimeout,
ReconcilerWorkers: config.ReconcilerWorkers,
LockClient: lockClient,
})
if err != nil {
return nil, fmt.Errorf("failed to initialize notification event handler: %w", err)
}
return eventHandler, nil
}
func NewNotificationService(
logger *slog.Logger,
adapter notification.Repository,
webhook notificationwebhook.Handler,
featureConnector feature.FeatureConnector,
) (notification.Service, error) {
notificationService, err := notificationservice.New(notificationservice.Config{
Adapter: adapter,
Webhook: webhook,
FeatureConnector: featureConnector,
Logger: logger.With(slog.String("subsystem", "notification")),
})
if err != nil {
return nil, fmt.Errorf("failed to initialize notification service: %w", err)
}
return notificationService, nil
}
func NewNoopNotificationWebhookHandler(
logger *slog.Logger,
) (notificationwebhook.Handler, error) {
return webhooknoop.New(logger), nil
}
func NewNotificationWebhookHandler(
logger *slog.Logger,
tracer trace.Tracer,
webhookConfig config.WebhookConfiguration,
svixClient *svix.Svix,
) (notificationwebhook.Handler, error) {
if svixClient == nil {
logger.Warn("svix client not configured, using noop handler")
return webhooknoop.New(logger), nil
}
handler, err := webhooksvix.New(webhooksvix.Config{
SvixAPIClient: svixClient,
RegisterEventTypes: notificationwebhook.NotificationEventTypes,
RegistrationTimeout: webhookConfig.EventTypeRegistrationTimeout,
SkipRegistrationOnError: webhookConfig.SkipEventTypeRegistrationOnError,
Logger: logger.WithGroup("notification.webhook"),
Tracer: tracer,
})
if err != nil {
return nil, fmt.Errorf("failed to initialize notification webhook handler: %w", err)
}
return handler, nil
}