Versions in this module Expand all Collapse all v0 v0.3.0 Sep 5, 2026 v0.2.0 Aug 21, 2026 Changes in this version + func SchemaSQL() string v0.1.1 Aug 5, 2026 v0.1.0 Jul 22, 2026 Changes in this version + const DefaultConversionTopic + const DefaultDLQCollection + const DefaultImpressionTopic + const DefaultTTL + const DefaultTokenCollection + const FCMMaxBatchSize + const FCMMaxPayloadSize + const TopicExperimentAssigned + const TopicNotificationFailed + const TopicNotificationSent + var ErrBackendUnavailable = errors.New("grnoti: backend unavailable") + var ErrCircuitOpen = errors.New("grnoti: circuit breaker is open") + var ErrClosed = errors.New("grnoti: closed") + var ErrDLQEventNotClaimed = errors.New("grnoti: dead-letter event is not in a claimed (retrying) state") + var ErrDLQEventNotFound = errors.New("grnoti: dead-letter event not found") + var ErrExperimentAlreadyExists = errors.New("grnoti: experiment already exists") + var ErrExperimentHasNoVariants = errors.New("grnoti: experiment has no variants") + var ErrExperimentNotFound = errors.New("grnoti: experiment not found") + var ErrFCMClientNil = errors.New("grnoti: fcm client is nil") + var ErrFCMPayloadTooLarge = errors.New("grnoti: fcm payload exceeds maximum size") + var ErrInvalidEventID = errors.New("grnoti: event id is required") + var ErrInvalidEventType = errors.New("grnoti: invalid event type") + var ErrInvalidPriority = errors.New("grnoti: invalid priority") + var ErrNoTargetSpecified = errors.New("grnoti: at least one of user id, anonymous id, or device tokens is required") + var ErrPreferencesNotFound = errors.New("grnoti: preferences not found") + var ErrPreferencesUserIDRequired = errors.New("grnoti: preferences user id is required") + var ErrTemplateNotFound = errors.New("grnoti: template not found for event type") + var ErrTooManyRequests = errors.New("grnoti: too many requests while circuit breaker is half-open") + var ErrWorkerPoolFull = errors.New("grnoti: worker pool queue is full") + var Version = "v0.1.0" + func FullJitterBackoff(base, max time.Duration, attempt int) time.Duration + func PublishAssigned(ctx context.Context, bus grevents.Bus, logger Logger, ...) + func PublishFailed(ctx context.Context, bus grevents.Bus, logger Logger, ...) + func PublishSent(ctx context.Context, bus grevents.Bus, logger Logger, ...) + type AnalyticsPublisher interface + Close func() error + PublishConversion func(ctx context.Context, userID, experimentID string) error + PublishImpression func(ctx context.Context, userID, experimentID, variantID string) error + func NewKafkaAnalyticsPublisher(cfg KafkaAnalyticsPublisherConfig) (AnalyticsPublisher, error) + type BatchSplitter interface + Deduplicate func(tokens []DeviceToken) []DeviceToken + Split func(tokens []DeviceToken, maxBatchSize int) [][]DeviceToken + func NewBatchSplitter() BatchSplitter + type CircuitBreaker interface + Execute func(ctx context.Context, fn func() error) error + GetStats func() CircuitBreakerStats + Reset func() + State func() CircuitState + func NewCircuitBreaker(maxFailures int, timeout, resetTimeout time.Duration) (CircuitBreaker, error) + func NewCircuitBreakerWithConfig(config CircuitBreakerConfig) (CircuitBreaker, error) + type CircuitBreakerConfig struct + Logger Logger + MaxFailures int + MaxHalfOpenRequests int + ResetTimeout time.Duration + Timeout time.Duration + type CircuitBreakerStats struct + ConsecutiveFailures int + LastFailureTime time.Time + LastStateChange time.Time + OpenedAt time.Time + State CircuitState + TimeUntilNextAttempt time.Duration + TotalFailures int64 + TotalRejections int64 + TotalSuccesses int64 + type CircuitState string + const CircuitStateClosed + const CircuitStateHalfOpen + const CircuitStateOpen + type DLQEvent struct + AttemptHistory []DLQRetryAttempt + CreatedAt time.Time + Event Event + EventID string + FailureReason string + FirstFailureAt time.Time + LastAttemptAt time.Time + MaxRetries int + NextRetryAt time.Time + RetryCount int + Status DLQStatus + UpdatedAt time.Time + type DLQHandler interface + ClaimRetryableEvents func(ctx context.Context, limit int) ([]*DLQEvent, error) + Close func() error + GetEventByID func(ctx context.Context, eventID string) (*DLQEvent, error) + MarkRetried func(ctx context.Context, eventID string, success bool, attemptErr error) error + PublishToDLQ func(ctx context.Context, event Event, failureReason string) error + PurgeExpiredEvents func(ctx context.Context, maxAge time.Duration) (int64, error) + func NewMemoryDLQHandler(maxRetries int, retryDelay, maxRetryDelay time.Duration) DLQHandler + func NewMongoDLQHandler(cfg MongoDLQHandlerConfig) (DLQHandler, error) + func NewPostgresDLQHandler(cfg PostgresDLQHandlerConfig) (DLQHandler, error) + type DLQRetryAttempt struct + AttemptNumber int + AttemptedAt time.Time + ErrorMessage string + Success bool + type DLQStatus string + const DLQStatusExhausted + const DLQStatusPending + const DLQStatusResolved + const DLQStatusRetrying + type DeviceToken struct + AnonymousID string + AppVersion string + CreatedAt time.Time + DeviceID string + IsActive bool + Platform Platform + Token string + UpdatedAt time.Time + UserID string + type DispatchResult struct + Errors []error + FailureByPlatform map[Platform]int + FailureCount int + InvalidTokens []string + RetryableErrors int + SuccessByPlatform map[Platform]int + SuccessCount int + func (d DispatchResult) HasFailures() bool + func (d DispatchResult) TotalCount() int + type Event struct + AnonymousID string + DeviceTokens []string + EventID string + ExperimentID string + Payload map[string]string + Priority Priority + Timestamp time.Time + Type EventType + UserID string + func (e Event) GetTargetID() string + func (e Event) HasDirectTokens() bool + func (e Event) IsAnonymous() bool + func (e Event) IsAuthenticated() bool + func (e Event) Validate() error + type EventConsumer interface + Close func() error + Start func(ctx context.Context, handler func(context.Context, Event) error) error + func NewKafkaEventConsumer(cfg KafkaConsumerConfig) (EventConsumer, error) + type EventType string + const EventTypeAccountVerification + const EventTypeCustom + const EventTypeGenericMarketing + const EventTypeGenericTransactional + const EventTypePasswordReset + const EventTypeSystemAlert + func (e EventType) IsValid() bool + func (e EventType) String() string + type EventTypeMetadata struct + CanBeScheduled bool + Category NotificationCategory + DefaultPriority Priority + Description string + MaxRetries int + RequiresImmediateDelivery bool + RetryDelayMultiplier float64 + ShouldIncludeInDigest bool + Transactional bool + type EventTypeRegistry interface + All func() []EventType + Lookup func(t EventType) (EventTypeMetadata, bool) + Register func(t EventType, meta EventTypeMetadata) error + func NewEventTypeRegistry() EventTypeRegistry + type Experiment struct + CreatedAt time.Time + Enabled bool + ID string + Name string + UpdatedAt time.Time + Variants []ExperimentVariant + type ExperimentAssignedPayload struct + ExperimentID string + Timestamp time.Time + UserID string + VariantID string + type ExperimentAssignment struct + AssignedAt time.Time + ExperimentID string + UserID string + VariantID string + type ExperimentEngine interface + AssignVariant func(ctx context.Context, userID string, experiment *Experiment) (*ExperimentVariant, error) + GetVariant func(ctx context.Context, userID string, experimentID string) (*ExperimentVariant, error) + TrackConversion func(ctx context.Context, userID string, experimentID string) error + TrackImpression func(ctx context.Context, userID string, experimentID string, variantID string) error + func NewCacheBackedExperimentEngine(cache grcache.Cache, analytics AnalyticsPublisher, bus grevents.Bus, ...) ExperimentEngine + func NewDeterministicExperimentEngine(analytics AnalyticsPublisher, bus grevents.Bus, logger Logger) ExperimentEngine + type ExperimentStore interface + Close func() error + CreateExperiment func(ctx context.Context, experiment *Experiment) error + DeleteExperiment func(ctx context.Context, experimentID string) error + GetExperiment func(ctx context.Context, experimentID string) (*Experiment, error) + ListExperiments func(ctx context.Context) ([]*Experiment, error) + UpdateExperiment func(ctx context.Context, experiment *Experiment) error + func NewMemoryExperimentStore() ExperimentStore + func NewPostgresExperimentStore(cfg PostgresConfig) (ExperimentStore, error) + type ExperimentVariant struct + ID string + Name string + Payload map[string]string + Weight int + type FCMClient interface + Send func(ctx context.Context, message *messaging.Message) (string, error) + SendEachForMulticast func(ctx context.Context, message *messaging.MulticastMessage) (*messaging.BatchResponse, error) + type FCMDispatcherConfig struct + EnableRetry bool + MaxRetryAttempts int + RetryBaseDelay time.Duration + RetryMaxDelay time.Duration + func DefaultFCMDispatcherConfig() FCMDispatcherConfig + type FCMDispatcherDeps struct + CircuitBreaker CircuitBreaker + Client FCMClient + Config FCMDispatcherConfig + Logger Logger + Metrics Metrics + RateLimiter RateLimiter + type FCMError struct + Code FCMErrorCode + Err error + Message string + Token string + func NewFCMError(code FCMErrorCode, token, message string, err error) *FCMError + func (e *FCMError) Error() string + func (e *FCMError) IsPermanent() bool + func (e *FCMError) IsRetryable() bool + func (e *FCMError) Unwrap() error + type FCMErrorCode string + const FCMErrorCodeInternal + const FCMErrorCodeInvalidArgument + const FCMErrorCodeQuotaExceeded + const FCMErrorCodeSenderIDMismatch + const FCMErrorCodeThirdPartyAuthErr + const FCMErrorCodeUnavailable + const FCMErrorCodeUnregistered + const FCMErrorCodeUnspecified + func (c FCMErrorCode) IsPermanent() bool + func (c FCMErrorCode) IsRetryable() bool + type IdempotencyRecord struct + EventID string + ExpiresAt time.Time + ProcessedAt time.Time + type IdempotencyStore interface + Close func() error + IsProcessed func(ctx context.Context, eventID string) (bool, error) + MarkProcessed func(ctx context.Context, eventID string, ttl time.Duration) error + func NewCacheIdempotencyStore(cache grcache.Cache) IdempotencyStore + type KafkaAnalyticsPublisherConfig struct + Brokers []string + ConversionTopic string + ImpressionTopic string + Logger Logger + SaramaConfig *sarama.Config + type KafkaConsumerConfig struct + Brokers []string + GroupID string + Logger Logger + SaramaConfig *sarama.Config + Topics []string + type LocaleResolver interface + GetDefaultLocale func() string + ResolveLocale func(ctx context.Context, userID string) (string, error) + ResolveLocaleForAnonymous func(ctx context.Context, anonymousID string) (string, error) + func NewPreferencesLocaleResolver(store PreferencesStore, fallbackLocale string) LocaleResolver + func NewStaticLocaleResolver(locale string) LocaleResolver + type LocalizationStore interface + GetLocalizedTemplate func(eventType EventType, locale string) (MessageTemplate, error) + GetSupportedLocales func(eventType EventType) []string + RegisterLocalizedTemplate func(eventType EventType, locale string, template MessageTemplate) error + func NewInMemoryLocalizationStore() LocalizationStore + type LocalizedTemplate struct + DefaultLocale string + Templates map[string]MessageTemplate + type Logger interface + Debug func(msg string, args ...any) + Error func(msg string, args ...any) + Info func(msg string, args ...any) + Warn func(msg string, args ...any) + func NopLogger() Logger + func OrNop(l Logger) Logger + type Message struct + Actions []NotificationAction + Badge *int + Body string + Category NotificationCategory + ChannelID string + CollapseKey string + Data map[string]string + DeepLink string + ImageURL string + Priority Priority + Sound string + TTL time.Duration + Title string + type MessageTemplate struct + Actions []NotificationAction + BodyTemplate string + Category NotificationCategory + ChannelID string + CollapseKey string + DeepLink string + DefaultData map[string]string + DefaultTTL time.Duration + Sound string + TitleTemplate string + type Metrics interface + IncEventsSkipped func(reason string) + IncInvalidTokens func(count int) + IncNotificationsFailed func(eventType EventType, platform Platform, count int) + IncNotificationsProcessed func() + IncNotificationsSent func(eventType EventType, platform Platform, count int) + ObserveDispatchLatency func(eventType EventType, platform Platform, duration time.Duration) + ObserveProcessingLatency func(duration time.Duration) + type MongoDLQHandlerConfig struct + CollectionName string + Database string + Logger Logger + MaxRetries int + MaxRetryDelay time.Duration + RetryDelay time.Duration + URI string + type MongoTokenStoreConfig struct + CollectionName string + Database string + Logger Logger + URI string + type NotificationAction struct + ID string + Icon string + Title string + URL string + type NotificationCategory string + const CategoryAlert + const CategoryMarketing + const CategorySocial + const CategoryTransactional + type NotificationFailedPayload struct + AnonymousID string + EventID string + EventType EventType + FailureCount int + Reason string + Timestamp time.Time + UserID string + type NotificationPreferences struct + CreatedAt time.Time + EventTypeSettings map[EventType]bool + GlobalEnabled bool + Locale string + QuietHoursEnabled bool + QuietHoursEnd string + QuietHoursStart string + Timezone string + UpdatedAt time.Time + UserID string + func (p *NotificationPreferences) IsEventTypeEnabled(eventType EventType) bool + type NotificationSentPayload struct + AnonymousID string + EventID string + EventType EventType + FailureCount int + SuccessCount int + Timestamp time.Time + UserID string + type NotificationService interface + Close func() error + ProcessEvent func(ctx context.Context, event Event) (ProcessingResult, error) + Submit func(ctx context.Context, event Event) error + func NewNotificationService(deps ServiceDeps) (NotificationService, error) + type NotificationTarget interface + GetTokens func() []DeviceToken + GetTopicName func() string + IsTopicBased func() bool + type PayloadValidator interface + EstimateSize func(msg Message) int + ValidateSize func(msg Message) error + func NewFCMPayloadValidator() PayloadValidator + type Platform string + const PlatformAndroid + const PlatformIOS + const PlatformWeb + func (p Platform) IsValid() bool + func (p Platform) String() string + type PostgresConfig struct + ConnectTimeout time.Duration + DSN string + Logger Logger + MaxConnLifetime time.Duration + MaxConns int32 + MinConns int32 + Pool *pgxpool.Pool + SkipSchemaEnsure bool + type PostgresDLQHandlerConfig struct + MaxRetries int + MaxRetryDelay time.Duration + RetryDelay time.Duration + type PreferencesFilter interface + ShouldSendNotification func(ctx context.Context, event Event) (bool, string, error) + func NewPreferencesFilter(store PreferencesStore, logger Logger) PreferencesFilter + type PreferencesStore interface + Close func() error + GetPreferences func(ctx context.Context, userID string) (*NotificationPreferences, error) + IsEventTypeEnabled func(ctx context.Context, userID string, eventType EventType) (bool, error) + SavePreferences func(ctx context.Context, prefs *NotificationPreferences) error + func NewCachedPreferencesStore(store PreferencesStore, cache grcache.Cache, ttl time.Duration, logger Logger) PreferencesStore + func NewMemoryPreferencesStore() PreferencesStore + func NewPostgresPreferencesStore(cfg PostgresConfig) (PreferencesStore, error) + type Priority string + const PriorityHigh + const PriorityLow + const PriorityNormal + func (p Priority) IsValid() bool + func (p Priority) String() string + type ProcessingResult struct + DispatchResult DispatchResult + Duration time.Duration + EventID string + ProcessedAt time.Time + SkipReason string + Skipped bool + TokenCount int + UserID string + type PushDispatcher interface + Send func(ctx context.Context, tokens []DeviceToken, msg Message) (DispatchResult, error) + SendToToken func(ctx context.Context, token DeviceToken, msg Message) error + SendToTopic func(ctx context.Context, topic string, msg Message) error + func NewFCMDispatcher(deps FCMDispatcherDeps) (PushDispatcher, error) + type RateLimiter interface + Allow func(ctx context.Context) (bool, error) + GetStats func(ctx context.Context) (RateLimiterStats, error) + Wait func(ctx context.Context) error + func NewLocalRateLimiter(requestsPerSecond, burstSize int) (RateLimiter, error) + func NewRedisRateLimiter(cfg RedisRateLimiterConfig) (RateLimiter, error) + type RateLimiterStats struct + AllowedCount int64 + BlockedCount int64 + BurstSize int + LastAllowedAt time.Time + RequestsPerSecond int + WaitCount int64 + type RedisRateLimiterConfig struct + Addr string + BurstSize int + DB int + DialTimeout time.Duration + Key string + Logger Logger + Password string + PoolSize int + ReadTimeout time.Duration + RequestsPerSecond int + WriteTimeout time.Duration + type RetryStrategy interface + GetDelay func(attempt int) time.Duration + ShouldRetry func(attempt int, err error) bool + func NewFullJitterRetry(maxAttempts int, baseDelay, maxDelay time.Duration) RetryStrategy + func NewNoopRetryStrategy() RetryStrategy + type ServiceConfig struct + EnableABTesting bool + EnableBackpressure bool + EnableDLQ bool + EnableEventBus bool + EnableLocalization bool + EnableMetrics bool + EnablePreferencesFilter bool + EnableRichPush bool + EnableTokenDeduplication bool + EnableTopicRouting bool + EnforceBatching bool + IdempotencyTTL time.Duration + MaxTokensPerBatch int + SkipInvalidEvents bool + func DefaultServiceConfig() ServiceConfig + type ServiceDeps struct + Config ServiceConfig + DLQHandler DLQHandler + Dispatcher PushDispatcher + EventBus grevents.Bus + Idempotency IdempotencyStore + Logger Logger + Metrics Metrics + PreferencesFilter PreferencesFilter + Templates TemplateEngine + TokenStore TokenStore + TopicRouter TopicRouter + WorkerPoolConfig WorkerPoolConfig + type TemplateEngine interface + BuildMessage func(event Event) (Message, error) + RegisterTemplate func(eventType EventType, template MessageTemplate) error + func NewLocalizedTemplateEngine(baseEngine TemplateEngine, localeStore LocalizationStore, ...) TemplateEngine + func NewTemplateEngine() TemplateEngine + type TokenStore interface + Close func() error + DeleteToken func(ctx context.Context, token string) error + GetActiveTokens func(ctx context.Context, userID string) ([]DeviceToken, error) + GetActiveTokensBatch func(ctx context.Context, userIDs []string) (map[string][]DeviceToken, error) + GetActiveTokensByAnonymousID func(ctx context.Context, anonymousID string) ([]DeviceToken, error) + MarkInvalid func(ctx context.Context, token string) error + SaveToken func(ctx context.Context, token DeviceToken) error + func NewMemoryTokenStore() TokenStore + func NewMongoTokenStore(cfg MongoTokenStoreConfig) (TokenStore, error) + func NewPostgresTokenStore(cfg PostgresConfig) (TokenStore, error) + type TokenTarget struct + Tokens []DeviceToken + func (t TokenTarget) GetTokens() []DeviceToken + func (t TokenTarget) GetTopicName() string + func (t TokenTarget) IsTopicBased() bool + type TopicRouter interface + ResolveTarget func(ctx context.Context, event Event) (NotificationTarget, error) + func NewEventTypeTopicRouter(topicMappings map[EventType]string, tokenStore TokenStore, logger Logger) TopicRouter + func NewStaticTopicRouter(topic string) TopicRouter + func NewTokenOnlyRouter(tokenStore TokenStore) TopicRouter + type TopicTarget struct + Topic string + func (t TopicTarget) GetTokens() []DeviceToken + func (t TopicTarget) GetTopicName() string + func (t TopicTarget) IsTopicBased() bool + type WorkerPool struct + func NewWorkerPool(deps WorkerPoolDeps) (*WorkerPool, error) + func (wp *WorkerPool) GetStats() WorkerPoolStats + func (wp *WorkerPool) Start() + func (wp *WorkerPool) Stop() + func (wp *WorkerPool) Submit(event Event) error + func (wp *WorkerPool) SubmitAsync(event Event) bool + type WorkerPoolConfig struct + QueueSize int + Workers int + type WorkerPoolDeps struct + Config WorkerPoolConfig + Handler func(context.Context, Event) error + Logger Logger + Metrics Metrics + type WorkerPoolStats struct + QueueSize int + QueueUsage float64 + QueuedEvents int + Workers int