Documentation
¶
Overview ¶
Package grnoti provides a push-notification service for the gourdian ecosystem: FCM dispatch, idempotent event processing, device-token management, durable dead-letter retry tracking, circuit breaking, distributed rate limiting, deterministic A/B experiment assignment, localization, and topic-based routing, behind a set of storage-agnostic interfaces.
Package shape ¶
grnoti's public API is a single flat package with no subpackages — every backend (MongoDB, PostgreSQL, Redis, Kafka, FCM) lives in this one module, distinguished by a "<concern>.<backend>.go" file-naming convention (e.g. tokenstore.mongo.go, dlq.postgres.go, ratelimiter.redis.go). The one exception is internal/postgresdb, sqlc's generated query code — it's a real Go subpackage, but an unexported internal/ one, not importable outside this module, so it doesn't undermine the "flat public API" claim. This is a deliberate divergence from sibling repos like grcache/graudit, which use one subpackage per backend to keep unused backend drivers out of a consumer's dependency graph — grnoti accepts that cost (importing grnoti pulls in the Mongo driver, pgx/v5 + sqlc-generated Postgres code, go-redis, sarama, and the Firebase messaging SDK regardless of which backends are actually used) in exchange for a simpler package structure, following gourdiantoken's precedent rather than grcache's. See docs/architecture.md for the full reasoning.
Precise, non-aspirational claims ¶
DLQHandler's durability guarantee is independent of grevents: an event dead-lettered here survives a process restart, unlike grevents' DeadLetterSink (an in-memory, best-effort recent-history buffer by design). grnoti optionally publishes lifecycle events ("notification.sent", "notification.failed", "experiment.assigned") via an injected grevents.Bus, but that publish is always best-effort — a nil bus or a publish failure never affects the durable operation it follows.
ExperimentEngine's variant assignment is a pure, deterministic function of (userID, experiment.ID, experiment.Variants): the same inputs always produce the same variant, with or without the optional assignment cache.
ServiceConfig.EnableRichPush, EnableLocalization, and EnableABTesting are composition-time bookkeeping flags only — nothing in notificationService.processEvent branches on them. What each would describe is decided by which concrete TemplateEngine/ExperimentEngine a caller composes into ServiceDeps before construction, not by a runtime check on these flags.
Index ¶
- Constants
- Variables
- 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, ...)
- func SchemaSQL() string
- type AnalyticsPublisher
- type BatchSplitter
- type CircuitBreaker
- type CircuitBreakerConfig
- type CircuitBreakerStats
- type CircuitState
- type DLQEvent
- type DLQHandler
- type DLQRetryAttempt
- type DLQStatus
- type DeviceToken
- type DispatchResult
- type Event
- type EventConsumer
- type EventType
- type EventTypeMetadata
- type EventTypeRegistry
- type Experiment
- type ExperimentAssignedPayload
- type ExperimentAssignment
- type ExperimentEngine
- type ExperimentStore
- type ExperimentVariant
- type FCMClient
- type FCMDispatcherConfig
- type FCMDispatcherDeps
- type FCMError
- type FCMErrorCode
- type IdempotencyRecord
- type IdempotencyStore
- type KafkaAnalyticsPublisherConfig
- type KafkaConsumerConfig
- type LocaleResolver
- type LocalizationStore
- type LocalizedTemplate
- type Logger
- type Message
- type MessageTemplate
- type Metrics
- type MongoDLQHandlerConfig
- type MongoTokenStoreConfig
- type NotificationAction
- type NotificationCategory
- type NotificationFailedPayload
- type NotificationPreferences
- type NotificationSentPayload
- type NotificationService
- type NotificationTarget
- type PayloadValidator
- type Platform
- type PostgresConfig
- type PostgresDLQHandlerConfig
- type PreferencesFilter
- type PreferencesStore
- type Priority
- type ProcessingResult
- type PushDispatcher
- type RateLimiter
- type RateLimiterStats
- type RedisRateLimiterConfig
- type RetryStrategy
- type ServiceConfig
- type ServiceDeps
- type TemplateEngine
- type TokenStore
- type TokenTarget
- type TopicRouter
- type TopicTarget
- type WorkerPool
- type WorkerPoolConfig
- type WorkerPoolDeps
- type WorkerPoolStats
Constants ¶
const ( // FCMMaxBatchSize is FCM's documented maximum tokens per multicast // request. FCMMaxBatchSize = 500 // DefaultTTL is applied to Android messages whose Message.TTL is unset. DefaultTTL = 24 * time.Hour )
const ( // TopicNotificationSent fires after a dispatch completes with at // least one successful delivery. TopicNotificationSent = "notification.sent" // TopicNotificationFailed fires after a dispatch ends with zero // successful deliveries. TopicNotificationFailed = "notification.failed" // TopicExperimentAssigned fires the first time a user is assigned a // variant for an experiment — not on subsequent lookups of an // already-assigned user. TopicExperimentAssigned = "experiment.assigned" )
Topic* are the grevents topics grnoti publishes to. Two dot-separated segments each, matching the real convention already established by grevents' own examples and graudit's TopicAuditRecorded (see docs/plan/grnoti-plan.md §1.2/§1.3).
const ( // DefaultImpressionTopic is used when // KafkaAnalyticsPublisherConfig.ImpressionTopic is empty. DefaultImpressionTopic = "grnoti.experiment.impressions" // DefaultConversionTopic is used when // KafkaAnalyticsPublisherConfig.ConversionTopic is empty. DefaultConversionTopic = "grnoti.experiment.conversions" )
const DefaultDLQCollection = "grnoti_dlq"
DefaultDLQCollection is the collection name used when MongoDLQHandlerConfig.CollectionName is empty.
const DefaultTokenCollection = "grnoti_tokens"
DefaultTokenCollection is the collection name used when MongoTokenStoreConfig.CollectionName is empty.
const FCMMaxPayloadSize = 4096
FCMMaxPayloadSize is FCM's documented maximum message payload size in bytes.
Variables ¶
var ( // ErrClosed indicates a method was called after Close. ErrClosed = errors.New("grnoti: closed") // reached (connection failure, timeout, transaction failure, etc.). ErrBackendUnavailable = errors.New("grnoti: backend unavailable") // ErrInvalidEventID indicates Event.EventID was empty. ErrInvalidEventID = errors.New("grnoti: event id is required") // ErrNoTargetSpecified indicates an Event has none of UserID, // AnonymousID, or DeviceTokens set — there is no way to resolve who // should receive it. ErrNoTargetSpecified = errors.New("grnoti: at least one of user id, anonymous id, or device tokens is required") // ErrInvalidEventType indicates Event.Type failed EventType.IsValid(). ErrInvalidEventType = errors.New("grnoti: invalid event type") // ErrInvalidPriority indicates Event.Priority failed Priority.IsValid(). ErrInvalidPriority = errors.New("grnoti: invalid priority") // ErrTemplateNotFound indicates no MessageTemplate is registered for // an event type, and no EventTypeCustom fallback is registered either. ErrTemplateNotFound = errors.New("grnoti: template not found for event type") // ErrPreferencesNotFound indicates no NotificationPreferences exist // for a user. Callers generally treat this as "use defaults," not a // hard failure — see PreferencesStore.IsEventTypeEnabled. ErrPreferencesNotFound = errors.New("grnoti: preferences not found") // ErrPreferencesUserIDRequired indicates PreferencesStore.SavePreferences // was called with an empty NotificationPreferences.UserID. Distinct from // ErrNoTargetSpecified, which is specifically about an Event having no // resolvable recipient — reusing it here would repeat the exact class of // sentinel-reuse-across-unrelated-conditions bug documented above. ErrPreferencesUserIDRequired = errors.New("grnoti: preferences user id is required") // ErrExperimentNotFound indicates no Experiment is registered under // the requested ID. ErrExperimentNotFound = errors.New("grnoti: experiment not found") // ErrExperimentAlreadyExists indicates ExperimentStore.CreateExperiment // was called with an ID that already exists. ErrExperimentAlreadyExists = errors.New("grnoti: experiment already exists") // ErrExperimentHasNoVariants indicates an Experiment exists but has // zero ExperimentVariants, so no assignment can be made. ErrExperimentHasNoVariants = errors.New("grnoti: experiment has no variants") // ErrDLQEventNotFound indicates a DLQHandler lookup found no // DLQEvent for the requested event ID. ErrDLQEventNotFound = errors.New("grnoti: dead-letter event not found") // ErrDLQEventNotClaimed indicates MarkRetried was called for an event // that is not currently in the "retrying" (claimed) state — either it // was never claimed via ClaimRetryableEvents, or a concurrent caller // already resolved/exhausted it. See docs/plan/grnoti-plan.md §5 for // why this replaces the source's unguarded read-then-write update. ErrDLQEventNotClaimed = errors.New("grnoti: dead-letter event is not in a claimed (retrying) state") // ErrFCMClientNil indicates a PushDispatcher was constructed with a // nil FCM client. ErrFCMClientNil = errors.New("grnoti: fcm client is nil") // ErrFCMPayloadTooLarge indicates a Message's estimated FCM payload // size exceeds FCMMaxPayloadSize. ErrFCMPayloadTooLarge = errors.New("grnoti: fcm payload exceeds maximum size") )
Sentinel errors for use with errors.Is. Backend implementations translate their own native errors (mongo.ErrNoDocuments, pgx.ErrNoRows, redis.Nil, ...) into these sentinels before wrapping with fmt.Errorf("...: %w", ...) — a backend-native error must never leak through a grnoti interface unwrapped, matching grcache's and graudit's own documented rule.
There is deliberately no IsX(err error) bool helper: callers use errors.Is(err, grnoti.ErrClosed) directly, consistent with every other gourdian repo's sentinel-error convention.
Two sentinels replace error reuse found in the reference implementation (see docs/plan/grnoti-plan.md §2 item 6): ErrNoTargetSpecified used to be ErrInvalidUserID reused for a semantically different condition, and ErrExperimentNotFound/ErrExperimentHasNoVariants used to both be ErrTemplateNotFound reused.
var ErrCircuitOpen = errors.New("grnoti: circuit breaker is open")
ErrCircuitOpen is returned by CircuitBreaker.Execute when the breaker is open and its Timeout has not yet elapsed.
var ErrTooManyRequests = errors.New("grnoti: too many requests while circuit breaker is half-open")
ErrTooManyRequests is returned by CircuitBreaker.Execute when the breaker is half-open and MaxHalfOpenRequests trial requests are already in flight.
var ErrWorkerPoolFull = errors.New("grnoti: worker pool queue is full")
ErrWorkerPoolFull is returned by WorkerPool.Submit/SubmitAsync when the queue is full.
var Version = "v0.2.0"
Version is the semantic version of this module, matching its most recent git tag.
Functions ¶
func FullJitterBackoff ¶
FullJitterBackoff returns a randomized backoff duration for the given 0-indexed attempt that just failed: sleep = random(0, min(cap, base*2^attempt)) — the AWS "Full Jitter" formula. This mirrors grevents' retry.go computeBackoff exactly (see docs/plan/grnoti-plan.md §1.2): grevents was the first backoff-with-jitter implementation in the gourdian ecosystem, and this is the second, deliberately kept identical rather than inventing a second formula. It is exported (unlike grevents' own unexported computeBackoff) so both retrystrategy.go's FCM-dispatch retry and the Postgres/Mongo DLQ backends' own retry-delay computation share one implementation instead of two independently-written copies of the same formula.
Parameters:
- base: time.Duration — the starting point; base<=0 returns 0 (no backoff)
- max: time.Duration — the ceiling; max<=0 defaults to defaultMaxBackoff
- attempt: int — 0-indexed attempt number
Returns:
- time.Duration: a value in [0, min(max, base*2^attempt)]
func PublishAssigned ¶
func PublishAssigned(ctx context.Context, bus grevents.Bus, logger Logger, payload ExperimentAssignedPayload)
PublishAssigned publishes a TopicExperimentAssigned event. See PublishSent for the nil-bus/best-effort contract.
Callers (deterministicExperimentEngine, cacheExperimentEngine) call this only on a genuinely new assignment, never on a lookup of an already-assigned user. The exactly-once-vs-at-least-once guarantee for a given (userID, experimentID) pair's first assignment differs by engine:
- deterministicExperimentEngine holds a single in-process lock across its whole check-then-write sequence, so exactly one goroutine can ever be the one that assigns+publishes for a given key — exactly-once (see experiment.go's own doc comment).
- cacheExperimentEngine cannot close the equivalent race, because grcache.Cache has no compare-and-swap/SetNX primitive to make its check-then-write atomic — concurrent racers on a brand-new pair can each independently publish, giving at-least-once delivery there (see cache.experiment.go's own doc comment).
Either way the underlying map/cache write itself stays correct (both racers compute the identical deterministic variant), and at-least-once is consistent with grevents' own Bus.Publish, which makes no exactly-once guarantee either.
func PublishFailed ¶
func PublishFailed(ctx context.Context, bus grevents.Bus, logger Logger, payload NotificationFailedPayload)
PublishFailed publishes a TopicNotificationFailed event. See PublishSent for the nil-bus/best-effort contract.
func PublishSent ¶
func PublishSent(ctx context.Context, bus grevents.Bus, logger Logger, payload NotificationSentPayload)
PublishSent publishes a TopicNotificationSent event for payload via bus. Following graudit's exact PublishRecorded precedent (docs/plan/ grnoti-plan.md §1.2/§1.3): bus may be nil (a silent no-op), and any error bus.Publish returns is only logged, never propagated to the caller — grevents delivery is a best-effort side channel on top of whatever durable/authoritative work already happened, never allowed to fail or block it.
func SchemaSQL ¶ added in v0.2.0
func SchemaSQL() string
SchemaSQL returns grnoti's Postgres schema (plain CREATE TABLE/INDEX IF NOT EXISTS, no grnoti-specific magic) as text, for the caller to apply through their own project's migration tool (golang-migrate, Flyway, a plain SQL file run in CI, whatever you already use) before constructing any of grnoti's four Postgres-backed stores:
- grnoti_tokens — NewPostgresTokenStore
- grnoti_preferences — NewPostgresPreferencesStore
- grnoti_experiments — NewPostgresExperimentStore
- grnoti_dlq — NewPostgresDLQHandler
grnoti deliberately never applies its own schema at runtime: doing so requires the runtime connection's role to have CREATE on the target schema, which a deliberately least-privilege application role (a common production setup — a separate role owns migrations, the app connects with a DML-only role) won't have. Vendor this text into your own migration once; re-sync it by hand (see CHANGELOG.md) if you upgrade grnoti and its schema changes. See docs/postgres.md for the full pattern, including how to share one pool across all four stores.
Types ¶
type AnalyticsPublisher ¶
type AnalyticsPublisher interface {
PublishImpression(ctx context.Context, userID, experimentID, variantID string) error
PublishConversion(ctx context.Context, userID, experimentID string) error
// Close releases any underlying connections. Idempotent.
Close() error
}
AnalyticsPublisher publishes experiment impression/conversion events for external analysis. New relative to the reference implementation, whose TrackImpression/TrackConversion were hardcoded no-ops — see docs/plan/grnoti-plan.md §2 item 9.
func NewKafkaAnalyticsPublisher ¶
func NewKafkaAnalyticsPublisher(cfg KafkaAnalyticsPublisherConfig) (AnalyticsPublisher, error)
NewKafkaAnalyticsPublisher builds a sarama.SyncProducer from cfg — synchronous rather than async so PublishImpression/PublishConversion's error return reflects an actual broker ack, not just a local enqueue.
type BatchSplitter ¶
type BatchSplitter interface {
// Split partitions tokens into batches of at most maxBatchSize each.
// maxBatchSize<=0 or an empty tokens returns tokens as one batch.
Split(tokens []DeviceToken, maxBatchSize int) [][]DeviceToken
// Deduplicate removes tokens with a repeated Token value, keeping the
// first occurrence (order-preserving).
Deduplicate(tokens []DeviceToken) []DeviceToken
}
BatchSplitter deduplicates and chunks DeviceTokens for dispatch.
func NewBatchSplitter ¶
func NewBatchSplitter() BatchSplitter
NewBatchSplitter returns a stateless BatchSplitter.
type CircuitBreaker ¶
type CircuitBreaker interface {
// Execute runs fn if the breaker's current state allows it.
//
// Returns:
// - error: ErrCircuitOpen if the breaker is open and its Timeout
// hasn't elapsed; ErrTooManyRequests if half-open and
// MaxHalfOpenRequests trial requests are already in flight;
// otherwise fn's own return value
Execute(ctx context.Context, fn func() error) error
State() CircuitState
GetStats() CircuitBreakerStats
// Reset forces the breaker back to CircuitStateClosed, for
// administrative use.
Reset()
}
CircuitBreaker wraps calls to an unreliable dependency (FCM) so persistent failures stop being retried immediately and instead fail fast for a cooldown period.
func NewCircuitBreaker ¶
func NewCircuitBreaker(maxFailures int, timeout, resetTimeout time.Duration) (CircuitBreaker, error)
NewCircuitBreaker constructs a CircuitBreaker with MaxHalfOpenRequests fixed at 1.
Parameters:
- maxFailures: int — consecutive failures before opening; must be > 0
- timeout: time.Duration — how long to stay open before allowing a trial request; must be > 0
- resetTimeout: time.Duration — how long a closed breaker must go without a failure before its consecutive-failure counter resets; must be > 0
Returns:
- CircuitBreaker
- error: non-nil if any parameter is not positive
func NewCircuitBreakerWithConfig ¶
func NewCircuitBreakerWithConfig(config CircuitBreakerConfig) (CircuitBreaker, error)
NewCircuitBreakerWithConfig constructs a CircuitBreaker from a full CircuitBreakerConfig.
Parameters:
- config: CircuitBreakerConfig — MaxFailures/Timeout/ResetTimeout must each be > 0; MaxHalfOpenRequests defaults to 1 if <= 0
Returns:
- CircuitBreaker
- error: non-nil if MaxFailures/Timeout/ResetTimeout is not positive
type CircuitBreakerConfig ¶
type CircuitBreakerConfig struct {
// MaxFailures is the number of consecutive failures that trips the
// breaker from closed to open.
MaxFailures int
// Timeout is how long the breaker stays open before allowing a trial
// request through (transitioning to half-open).
Timeout time.Duration
// ResetTimeout is how long a closed breaker must go without a failure
// before its consecutive-failure counter resets to zero.
ResetTimeout time.Duration
// MaxHalfOpenRequests bounds concurrent trial requests while
// half-open. Defaults to 1 if <= 0.
MaxHalfOpenRequests int
// Logger receives optional diagnostic messages for state transitions
// (open/half-open/close). A nil Logger disables logging.
Logger Logger
}
CircuitBreakerConfig configures a CircuitBreaker.
type CircuitBreakerStats ¶
type CircuitBreakerStats struct {
State CircuitState
ConsecutiveFailures int
TotalSuccesses int64
TotalFailures int64
TotalRejections int64
LastFailureTime time.Time
LastStateChange time.Time
OpenedAt time.Time
TimeUntilNextAttempt time.Duration
}
CircuitBreakerStats is a point-in-time snapshot of a CircuitBreaker's counters.
type CircuitState ¶
type CircuitState string
CircuitState is a CircuitBreaker's current state.
const ( // CircuitStateClosed is the normal state: requests pass through, and // consecutive failures are counted toward MaxFailures. CircuitStateClosed CircuitState = "closed" // CircuitStateOpen rejects every request immediately (without // attempting them) until Timeout elapses, then transitions to // CircuitStateHalfOpen. CircuitStateOpen CircuitState = "open" // CircuitStateHalfOpen allows up to MaxHalfOpenRequests trial requests // through to test whether the failing dependency has recovered: the // first trial failure trips it straight back to CircuitStateOpen, the // first trial success closes it. CircuitStateHalfOpen CircuitState = "half_open" )
type DLQEvent ¶
type DLQEvent struct {
EventID string
Event Event
FailureReason string
RetryCount int
MaxRetries int
FirstFailureAt time.Time
LastAttemptAt time.Time
NextRetryAt time.Time
Status DLQStatus
AttemptHistory []DLQRetryAttempt
CreatedAt time.Time
UpdatedAt time.Time
}
DLQEvent is one durably-tracked failed-delivery record.
type DLQHandler ¶
type DLQHandler interface {
// PublishToDLQ records a new failure for event, or (if event.EventID
// already has a pending/retrying record) appends to its existing
// attempt history.
PublishToDLQ(ctx context.Context, event Event, failureReason string) error
// ClaimRetryableEvents atomically selects up to limit events whose
// NextRetryAt has passed and whose Status is DLQStatusPending, and
// transitions each claimed event to DLQStatusRetrying as part of the
// same operation — so that N concurrent callers (e.g. N worker
// replicas) each claim disjoint events, never the same one twice.
// This replaces the reference implementation's GetRetryableEvents,
// which was a plain read with no claim semantics at all (see
// docs/plan/grnoti-plan.md §1.3, §2 item 4). Every event returned here
// must eventually be resolved via MarkRetried — an event claimed but
// never marked stays in DLQStatusRetrying until a backend-specific
// claim-timeout sweep (if configured) reclaims it.
//
// On error, the returned slice is not necessarily empty and callers
// must still process it: some backends (e.g. the Mongo implementation,
// which claims one document per iteration rather than in a single
// atomic statement) can already have durably transitioned a prefix of
// events to DLQStatusRetrying before hitting a failure on a later one.
// Discarding a non-nil slice just because err != nil would orphan
// those already-claimed events — there is no reclaim-timeout sweep in
// this package to recover them otherwise. Backends whose claim is a
// single atomic statement (e.g. Postgres) are all-or-nothing by
// necessity and always return a nil slice on error; that is a
// backend-specific limitation, not the general contract.
ClaimRetryableEvents(ctx context.Context, limit int) ([]*DLQEvent, error)
// MarkRetried records the outcome of a retry attempt for eventID and
// transitions it out of DLQStatusRetrying: to DLQStatusResolved on
// success, DLQStatusExhausted if retries are exhausted, or back to
// DLQStatusPending (with a recomputed NextRetryAt) otherwise.
//
// Returns:
// - error: wraps ErrDLQEventNotClaimed if eventID is not currently
// DLQStatusRetrying (already resolved by a concurrent caller, or
// never claimed) — implementations must scope their update to the
// claimed state rather than unconditionally overwriting, see
// docs/plan/grnoti-plan.md §5
MarkRetried(ctx context.Context, eventID string, success bool, attemptErr error) error
// GetEventByID returns a specific DLQEvent by ID regardless of status.
//
// Returns:
// - error: wraps ErrDLQEventNotFound if no such event exists
GetEventByID(ctx context.Context, eventID string) (*DLQEvent, error)
// PurgeExpiredEvents deletes DLQStatusResolved/DLQStatusExhausted
// events, and any event older than maxAge regardless of status.
//
// Returns:
// - int64: number of events deleted
PurgeExpiredEvents(ctx context.Context, maxAge time.Duration) (int64, error)
// Close releases any underlying connections. Idempotent.
Close() error
}
DLQHandler durably tracks push-delivery failures across retries and process restarts — a different durability contract than grevents' own in-memory DeadLetterSink, see docs/plan/grnoti-plan.md §1.2.
func NewMemoryDLQHandler ¶
func NewMemoryDLQHandler(maxRetries int, retryDelay, maxRetryDelay time.Duration) DLQHandler
NewMemoryDLQHandler constructs an in-memory DLQHandler.
Parameters:
- maxRetries: int — defaults to 3 if <= 0
- retryDelay, maxRetryDelay: time.Duration — passed to FullJitterBackoff for computing each event's NextRetryAt; unlike maxRetries, 0 is a valid, deliberate choice here (immediate retry-eligibility, useful for tests), not silently replaced with a default — pass 5*time.Minute/time.Hour explicitly for the values the Postgres/Mongo backends use as their own defaults
func NewMongoDLQHandler ¶
func NewMongoDLQHandler(cfg MongoDLQHandlerConfig) (DLQHandler, error)
NewMongoDLQHandler connects to MongoDB per cfg, ensures indexes (including a 7-day TTL index on created_at as a durable-retention backstop independent of PurgeExpiredEvents), and validates connectivity before returning.
Claim semantics (see docs/plan/grnoti-plan.md §1.3, §5): unlike the reference implementation's MongoDLQHandler.MarkRetried, which read RetryCount then wrote it back with no guard at all (a confirmed lost-update race under concurrent retries), every write here is scoped by an atomic MongoDB operation — ClaimRetryableEvents uses FindOneAndUpdate per document (atomic per-document claim, no transaction needed), and MarkRetried's retry_count increment is a $inc scoped to {event_id, status: "retrying"} rather than a Go-side read-then-set.
func NewPostgresDLQHandler ¶
func NewPostgresDLQHandler(cfg PostgresDLQHandlerConfig) (DLQHandler, error)
NewPostgresDLQHandler connects per cfg.
Claim semantics: ClaimRetryableEvents runs a single UPDATE statement whose subquery uses `SELECT ... FOR UPDATE SKIP LOCKED` to let N concurrent callers each claim a disjoint batch of pending events without contention — deliberately not graudit's pg_advisory_xact_lock (a single global serialization point, correct for graudit's one hash chain but wrong here, where claiming should be embarrassingly parallel across worker replicas). See docs/plan/grnoti-plan.md §1.3 and internal/postgresdb/queries/dlq.sql.
type DLQRetryAttempt ¶
type DLQRetryAttempt struct {
AttemptNumber int
AttemptedAt time.Time
Success bool
ErrorMessage string
}
DLQRetryAttempt records the outcome of one retry attempt for a DLQEvent.
type DLQStatus ¶
type DLQStatus string
DLQStatus is a DLQEvent's lifecycle state.
const ( // DLQStatusPending is newly-recorded or awaiting its next retry. DLQStatusPending DLQStatus = "pending" // DLQStatusRetrying means a worker currently holds an atomic claim on // this event (see DLQHandler.ClaimRetryableEvents) and is attempting // delivery — not a status any caller sets directly. DLQStatusRetrying DLQStatus = "retrying" // DLQStatusExhausted means retries were exhausted without success. DLQStatusExhausted DLQStatus = "exhausted" // DLQStatusResolved means a retry eventually succeeded. DLQStatusResolved DLQStatus = "resolved" )
type DeviceToken ¶
type DeviceToken struct {
Token string `json:"token" bson:"token"`
Platform Platform `json:"platform" bson:"platform"`
UserID string `json:"user_id,omitempty" bson:"user_id,omitempty"`
AnonymousID string `json:"anonymous_id,omitempty" bson:"anonymous_id,omitempty"`
DeviceID string `json:"device_id,omitempty" bson:"device_id,omitempty"`
AppVersion string `json:"app_version,omitempty" bson:"app_version,omitempty"`
IsActive bool `json:"is_active" bson:"is_active"`
CreatedAt time.Time `json:"created_at" bson:"created_at"`
UpdatedAt time.Time `json:"updated_at" bson:"updated_at"`
}
DeviceToken is a single registered push-notification destination.
type DispatchResult ¶
type DispatchResult struct {
// SuccessCount is how many recipients (tokens or, for a topic send, the
// single topic) were accepted for delivery.
SuccessCount int
// FailureCount is how many recipients failed, permanently or
// retryably — always >= len(InvalidTokens), since every invalid token
// is also counted as a failure.
FailureCount int
// InvalidTokens lists tokens FCM reported as permanently dead
// (unregistered/sender-mismatch) — the caller is expected to pass
// these to TokenStore.MarkInvalid rather than retry them.
InvalidTokens []string
// RetryableErrors counts failures classified as transient (see
// FCMErrorCode.IsRetryable) — a subset of FailureCount worth another
// attempt, as opposed to a permanent per-token failure already
// captured in InvalidTokens.
RetryableErrors int
// Errors collects the individual per-batch/per-token errors behind
// FailureCount, for logging — see joinDispatchErrors for how
// NotificationService summarizes these into one DLQ failure reason.
Errors []error
// SuccessByPlatform/FailureByPlatform break SuccessCount/FailureCount
// down per Platform — populated by dispatcher.fcm.go's Send, which
// dispatches each platform group separately before merging into the
// aggregate counts above. NotificationService uses this breakdown to
// call Metrics.IncNotificationsSent/Failed and
// Metrics.ObserveDispatchLatency with a real per-call Platform label
// (required by the Metrics interface) instead of guessing one.
SuccessByPlatform map[Platform]int
FailureByPlatform map[Platform]int
}
DispatchResult summarizes the outcome of a PushDispatcher.Send call.
func (DispatchResult) HasFailures ¶
func (d DispatchResult) HasFailures() bool
HasFailures reports whether FailureCount > 0.
func (DispatchResult) TotalCount ¶
func (d DispatchResult) TotalCount() int
TotalCount returns SuccessCount + FailureCount.
type Event ¶
type Event struct {
EventID string `json:"event_id" bson:"event_id"`
UserID string `json:"user_id,omitempty" bson:"user_id,omitempty"`
AnonymousID string `json:"anonymous_id,omitempty" bson:"anonymous_id,omitempty"`
DeviceTokens []string `json:"device_tokens,omitempty" bson:"device_tokens,omitempty"`
Type EventType `json:"type" bson:"type"`
Payload map[string]string `json:"payload" bson:"payload"`
Priority Priority `json:"priority" bson:"priority"`
Timestamp time.Time `json:"timestamp" bson:"timestamp"`
ExperimentID string `json:"experiment_id,omitempty" bson:"experiment_id,omitempty"`
}
Event is the unit of work a NotificationService processes: a single notification-worthy occurrence for one target (an authenticated user, an anonymous visitor, or a fixed set of device tokens).
func (Event) GetTargetID ¶
GetTargetID returns whichever identifier best names e's recipient, for logging/metrics — UserID, else AnonymousID, else the literal "direct" for a token-only Event.
func (Event) HasDirectTokens ¶
HasDirectTokens reports whether e targets an explicit set of device tokens, bypassing TokenStore lookup entirely.
func (Event) IsAnonymous ¶
IsAnonymous reports whether e targets an anonymous visitor (AnonymousID set, UserID not).
func (Event) IsAuthenticated ¶
IsAuthenticated reports whether e targets a known user (UserID set).
func (Event) Validate ¶
Validate checks that e has the minimum fields required to be processed.
Returns:
- error: ErrInvalidEventID, ErrNoTargetSpecified, ErrInvalidEventType, or ErrInvalidPriority — each a distinct sentinel for a distinct condition (the reference implementation reused one sentinel for two of these; see docs/plan/grnoti-plan.md §2 item 6)
type EventConsumer ¶
type EventConsumer interface {
// Start begins consuming and invoking handler for each Event, blocking
// until ctx is canceled or an unrecoverable error occurs.
Start(ctx context.Context, handler func(context.Context, Event) error) error
// Close stops consuming and releases underlying connections.
// Idempotent.
Close() error
}
EventConsumer ingests Events from an external source (Kafka) and invokes handler for each.
func NewKafkaEventConsumer ¶
func NewKafkaEventConsumer(cfg KafkaConsumerConfig) (EventConsumer, error)
NewKafkaEventConsumer creates a Kafka-backed EventConsumer, validating broker connectivity by constructing the underlying consumer group (which sarama connects eagerly, unlike a lazy client).
type EventType ¶
type EventType string
EventType identifies what kind of notification an Event represents. It is a plain string type, not a closed/sealed enum — grnoti ships a small, domain-neutral vocabulary (below) plus EventTypeCustom as an escape hatch, and consumers register their own application-specific types (and optional metadata) via EventTypeRegistry rather than grnoti maintaining an exhaustive catalog.
This is a deliberate departure from the reference implementation, which compiled ~130 e-commerce-specific constants directly into the library and spread each type's behavioral traits (default priority, category, retryability, ...) across eight separately-maintained exhaustive switch statements that had to be extended in lockstep for every new type. Here, a trait is one field in one EventTypeMetadata value, registered once.
const ( // EventTypeCustom is the fallback type for anything not otherwise // registered. TemplateEngine implementations fall back to a template // registered under this type when an event's exact Type has none of // its own. EventTypeCustom EventType = "custom" // EventTypeSystemAlert is a generic operational/system notification // (e.g. a security alert, a service disruption notice). EventTypeSystemAlert EventType = "system_alert" // EventTypeAccountVerification is a generic "verify your account" // notification. EventTypeAccountVerification EventType = "account_verification" // EventTypePasswordReset is a generic "reset your password" // notification. EventTypePasswordReset EventType = "password_reset" // EventTypeGenericTransactional is a generic transactional // notification with no more specific type registered — expected to be // delivered promptly and not batched into a digest. EventTypeGenericTransactional EventType = "generic_transactional" // EventTypeGenericMarketing is a generic promotional/marketing // notification — expected to respect quiet hours and be eligible for // digesting, unlike transactional types. EventTypeGenericMarketing EventType = "generic_marketing" )
A small, generic starter vocabulary. Deliberately not domain-specific — contrast with the reference implementation's ~130 e-commerce constants (order lifecycle, returns, loyalty, EMI due dates, ...), which belong in a consumer-side package, not here.
func (EventType) IsValid ¶
IsValid reports whether e is structurally usable as an event type — a non-empty string. This is intentionally not "is this a known/registered type": an Event carrying an application-specific EventType that was never registered with an EventTypeRegistry is still a structurally valid Event. Registration is opt-in metadata for consumers who want priority/category/retry defaults driven by event type, not a gate on which types may be used at all.
type EventTypeMetadata ¶
type EventTypeMetadata struct {
// DefaultPriority is used when an Event of this type doesn't specify
// its own Priority.
DefaultPriority Priority
// Category classifies the type for consumers that branch on it (e.g.
// PreferencesFilter's per-category opt-out).
Category NotificationCategory
// Transactional marks a type as transactional (order confirmations,
// security alerts) as opposed to marketing.
Transactional bool
// RequiresImmediateDelivery marks a type that should never be delayed
// or batched.
RequiresImmediateDelivery bool
// CanBeScheduled marks a type eligible for a future-scheduled send.
CanBeScheduled bool
// ShouldIncludeInDigest marks a type eligible for batching into a
// periodic digest notification instead of sending immediately.
ShouldIncludeInDigest bool
// MaxRetries and RetryDelayMultiplier tune dispatch retry behavior for
// this type; a zero MaxRetries means "use the dispatcher's default."
MaxRetries int
RetryDelayMultiplier float64
// Description is a short human-readable description, useful for admin
// tooling listing registered types.
Description string
}
EventTypeMetadata describes the behavioral traits associated with an EventType, registered via EventTypeRegistry.Register. Replaces the reference implementation's eight separate exhaustive switch statements (one per trait) with a single data value per type.
type EventTypeRegistry ¶
type EventTypeRegistry interface {
// Register associates meta with t, overwriting any existing
// registration for t.
//
// Parameters:
// - t: EventType — must satisfy IsValid()
// - meta: EventTypeMetadata
//
// Returns:
// - error: non-nil if t is not IsValid()
Register(t EventType, meta EventTypeMetadata) error
// Lookup returns the registered metadata for t, if any.
//
// Parameters:
// - t: EventType
//
// Returns:
// - EventTypeMetadata: the zero value if not found
// - bool: true if t has a registration
Lookup(t EventType) (EventTypeMetadata, bool)
// All returns every currently-registered EventType, in no particular
// order.
All() []EventType
}
EventTypeRegistry tracks EventTypeMetadata for known EventTypes. Consumer applications register their own types (and grnoti's own small starter vocabulary is pre-registered by NewEventTypeRegistry) instead of grnoti maintaining an exhaustive built-in catalog.
func NewEventTypeRegistry ¶
func NewEventTypeRegistry() EventTypeRegistry
NewEventTypeRegistry constructs an EventTypeRegistry pre-seeded with grnoti's own small generic vocabulary (EventTypeSystemAlert, EventTypeAccountVerification, EventTypePasswordReset, EventTypeGenericTransactional, EventTypeGenericMarketing) plus EventTypeCustom. Consumers call Register to add their own types on top.
type Experiment ¶
type Experiment struct {
ID string
Name string
Variants []ExperimentVariant
Enabled bool
CreatedAt time.Time
UpdatedAt time.Time
}
Experiment is an A/B (or A/B/n) test definition.
type ExperimentAssignedPayload ¶
type ExperimentAssignedPayload struct {
UserID string
ExperimentID string
VariantID string
Timestamp time.Time
}
ExperimentAssignedPayload is published on TopicExperimentAssigned.
type ExperimentAssignment ¶
type ExperimentAssignment struct {
UserID string
ExperimentID string
VariantID string
AssignedAt time.Time
}
ExperimentAssignment is one user's deterministic variant assignment for one Experiment.
type ExperimentEngine ¶
type ExperimentEngine interface {
// GetVariant returns userID's existing assignment for experimentID, if
// one has already been made and cached.
//
// Returns:
// - *ExperimentVariant: nil if no assignment exists yet — this is not
// an error; the caller should call AssignVariant to create one
// - error: only for a genuine operational failure reading the cache
GetVariant(ctx context.Context, userID string, experimentID string) (*ExperimentVariant, error)
// AssignVariant deterministically computes (and caches) userID's
// variant within experiment: the same (userID, experiment.ID,
// experiment.Variants) always produces the same variant, so repeated
// calls are stable even without the cache. On a genuinely new
// assignment, implementations publish a TopicExperimentAssigned event
// via PublishAssigned — see that function's doc comment for how the
// exactly-once-vs-at-least-once guarantee for that publish differs
// between this package's two implementations.
//
// Returns:
// - error: wraps ErrExperimentHasNoVariants if experiment.Variants is
// empty
AssignVariant(ctx context.Context, userID string, experiment *Experiment) (*ExperimentVariant, error)
// TrackImpression records that userID was shown variantID of
// experimentID, publishing a real analytics event (see
// AnalyticsPublisher) — unlike the reference implementation, this is
// not a no-op.
TrackImpression(ctx context.Context, userID string, experimentID string, variantID string) error
// TrackConversion records that userID converted while assigned to
// experimentID.
TrackConversion(ctx context.Context, userID string, experimentID string) error
}
ExperimentEngine computes deterministic variant assignment and records impression/conversion analytics. Unlike the reference implementation, this interface takes an *Experiment as input to AssignVariant rather than owning a mutable map of experiment definitions itself — assignment is a pure function of (userID, experiment), so a correct implementation needs no internal synchronization for the assignment computation itself; any caching of the result (see NewDeterministicExperimentEngine and NewCacheBackedExperimentEngine) is a separate, explicitly-synchronized concern.
func NewCacheBackedExperimentEngine ¶
func NewCacheBackedExperimentEngine(cache grcache.Cache, analytics AnalyticsPublisher, bus grevents.Bus, logger Logger) ExperimentEngine
NewCacheBackedExperimentEngine constructs an ExperimentEngine whose assignment cache is any grcache.Cache, for multi-instance deployments.
Parameters:
- cache: grcache.Cache — caller-owned; not closed by this engine (it has no Close method at all — see the ExperimentEngine interface)
- analytics: AnalyticsPublisher — may be nil, see TrackImpression
- bus: grevents.Bus — may be nil; AssignVariant publishes TopicExperimentAssigned on a new assignment when set (§1.2)
- logger: Logger — may be nil
func NewDeterministicExperimentEngine ¶
func NewDeterministicExperimentEngine(analytics AnalyticsPublisher, bus grevents.Bus, logger Logger) ExperimentEngine
NewDeterministicExperimentEngine constructs an ExperimentEngine with an in-process assignment cache. See cache.experiment.go (Stage 4) for a grcache-backed variant suited to multi-instance deployments, where the in-process cache here would give each instance its own (still individually-correct, since assignment is deterministic) cache instead of a shared one.
Parameters:
- analytics: AnalyticsPublisher — may be nil; TrackImpression/ TrackConversion log and no-op rather than erroring when unset
- bus: grevents.Bus — may be nil; AssignVariant publishes TopicExperimentAssigned on a new assignment when set (§1.2)
- logger: Logger — may be nil
type ExperimentStore ¶
type ExperimentStore interface {
CreateExperiment(ctx context.Context, experiment *Experiment) error
// GetExperiment returns experimentID's definition.
//
// Returns:
// - error: wraps ErrExperimentNotFound if no such experiment exists
GetExperiment(ctx context.Context, experimentID string) (*Experiment, error)
UpdateExperiment(ctx context.Context, experiment *Experiment) error
// DeleteExperiment removes an experiment definition. Deleting a
// nonexistent experiment is not an error.
DeleteExperiment(ctx context.Context, experimentID string) error
ListExperiments(ctx context.Context) ([]*Experiment, error)
// Close releases any underlying connections. Idempotent.
Close() error
}
ExperimentStore persists Experiment definitions. Assignment (which variant a given user gets) is deliberately not part of this interface — see ExperimentEngine, which computes assignment as a pure function of an Experiment fetched from here, rather than this store owning mutable per-user assignment state itself.
func NewMemoryExperimentStore ¶
func NewMemoryExperimentStore() ExperimentStore
NewMemoryExperimentStore constructs an in-memory ExperimentStore.
func NewPostgresExperimentStore ¶
func NewPostgresExperimentStore(cfg PostgresConfig) (ExperimentStore, error)
NewPostgresExperimentStore connects per cfg — CRUD storage for Experiment definitions (see docs/plan/grnoti-plan.md §6; the deterministic assignment algorithm itself lives in experiment.go/ cache.experiment.go, not here).
type ExperimentVariant ¶
type ExperimentVariant struct {
ID string
Name string
Weight int // relative weight for deterministic bucketing; weights need not sum to 100
Payload map[string]string
}
ExperimentVariant is one arm of an Experiment.
type FCMClient ¶
type FCMClient interface {
SendEachForMulticast(ctx context.Context, message *messaging.MulticastMessage) (*messaging.BatchResponse, error)
Send(ctx context.Context, message *messaging.Message) (string, error)
}
FCMClient is the subset of Firebase Admin SDK's *messaging.Client that fcmDispatcher needs. Unlike every other backend in this repo, which is tested against a real local instance (Mongo/Postgres/Redis/Kafka all run in docker for their own test suites), FCM has no local emulator for actually delivering pushes — this interface exists so fcmDispatcher's own logic (batching, retry, error classification, rate-limiter/ circuit-breaker wiring) can be tested against a fake, the one deliberate exception to this repo's real-services testing policy. Matches the reference implementation's own justification (fcm.dispatcher.go:26-27).
type FCMDispatcherConfig ¶
type FCMDispatcherConfig struct {
// EnableRetry turns on sendBatchWithRetry. If false, a batch is sent
// exactly once.
EnableRetry bool
// MaxRetryAttempts is the total number of attempts per batch,
// including the first. Required to be > 0 when EnableRetry is true —
// unlike a retry *delay*, a retry *count* of 0 has no sane meaning
// ("retry enabled but never send") to silently default around.
MaxRetryAttempts int
// RetryBaseDelay/RetryMaxDelay feed FullJitterBackoff between
// attempts.
RetryBaseDelay time.Duration
RetryMaxDelay time.Duration
}
FCMDispatcherConfig holds fcmDispatcher's retry tuning. See DefaultFCMDispatcherConfig for the recommended starting point.
func DefaultFCMDispatcherConfig ¶
func DefaultFCMDispatcherConfig() FCMDispatcherConfig
DefaultFCMDispatcherConfig returns a sane starting configuration: retry enabled, 3 attempts, 500ms base / 5s max jittered backoff.
type FCMDispatcherDeps ¶
type FCMDispatcherDeps struct {
// Client is required.
Client FCMClient
Config FCMDispatcherConfig
// RateLimiter, if set, gates every outbound batch/single-send through
// Wait before it reaches Client — the reference implementation built a
// RateLimiter but never connected it to dispatch at all
// (docs/plan/grnoti-plan.md §3.2).
RateLimiter RateLimiter
// CircuitBreaker, if set, wraps every Client call — same §3.2 gap for
// CircuitBreaker.
CircuitBreaker CircuitBreaker
// Metrics, if set, receives IncInvalidTokens for tokens FCM reports as
// permanently invalid. Send's per-eventType/per-platform metrics
// (IncNotificationsSent/Failed, ObserveDispatchLatency) are NOT called
// from here: PushDispatcher.Send only ever sees tokens+Message, never
// the originating Event/EventType those calls require — that wiring
// belongs one layer up, in NotificationService (Stage 12), which does
// have the Event.
Metrics Metrics
Logger Logger
}
FCMDispatcherDeps configures an fcmDispatcher constructed by NewFCMDispatcher, following WorkerPoolDeps' shape (required collaborator + tuning Config + optional collaborators + Logger).
type FCMError ¶
type FCMError struct {
Code FCMErrorCode
Token string
Message string
Err error
}
FCMError wraps a single FCM send failure for one token with its classified FCMErrorCode.
func NewFCMError ¶
func NewFCMError(code FCMErrorCode, token, message string, err error) *FCMError
NewFCMError constructs an FCMError.
func (*FCMError) IsPermanent ¶
IsPermanent delegates to e.Code.IsPermanent.
func (*FCMError) IsRetryable ¶
IsRetryable delegates to e.Code.IsRetryable.
type FCMErrorCode ¶
type FCMErrorCode string
FCMErrorCode classifies an FCM send failure into a small set of well-known categories, used to decide retryability without callers needing to string-match raw FCM SDK error text themselves.
const ( FCMErrorCodeUnspecified FCMErrorCode = "unspecified" FCMErrorCodeInvalidArgument FCMErrorCode = "invalid_argument" FCMErrorCodeUnregistered FCMErrorCode = "unregistered" FCMErrorCodeSenderIDMismatch FCMErrorCode = "sender_id_mismatch" FCMErrorCodeQuotaExceeded FCMErrorCode = "quota_exceeded" FCMErrorCodeInternal FCMErrorCode = "internal" FCMErrorCodeThirdPartyAuthErr FCMErrorCode = "third_party_auth_error" )
func (FCMErrorCode) IsPermanent ¶
func (c FCMErrorCode) IsPermanent() bool
IsPermanent reports whether a failure of this class means the token itself is dead and should be removed via TokenStore.MarkInvalid, rather than retried.
func (FCMErrorCode) IsRetryable ¶
func (c FCMErrorCode) IsRetryable() bool
IsRetryable reports whether a failure of this class is worth retrying (transient server-side conditions), as opposed to a condition that will never succeed no matter how many times it's retried.
type IdempotencyRecord ¶
type IdempotencyRecord struct {
EventID string `json:"event_id" bson:"event_id"`
ProcessedAt time.Time `json:"processed_at" bson:"processed_at"`
ExpiresAt time.Time `json:"expires_at" bson:"expires_at"`
}
IdempotencyRecord is what an IdempotencyStore persists for a processed event.
type IdempotencyStore ¶
type IdempotencyStore interface {
// IsProcessed reports whether eventID has already been marked
// processed and that mark has not yet expired.
IsProcessed(ctx context.Context, eventID string) (bool, error)
// MarkProcessed records eventID as processed for ttl. ttl<=0 means no
// expiry. Calling MarkProcessed twice for the same eventID is not an
// error — implementations treat it as idempotent by design.
MarkProcessed(ctx context.Context, eventID string, ttl time.Duration) error
// Close releases any underlying connections. Idempotent.
Close() error
}
IdempotencyStore records which events have already been processed, so a redelivered Event (e.g. from Kafka consumer-group rebalance) is not dispatched twice.
func NewCacheIdempotencyStore ¶
func NewCacheIdempotencyStore(cache grcache.Cache) IdempotencyStore
NewCacheIdempotencyStore constructs an IdempotencyStore backed by cache.
Parameters:
- cache: grcache.Cache — caller-owned; not closed by this store's Close (see Close's doc comment)
type KafkaAnalyticsPublisherConfig ¶
type KafkaAnalyticsPublisherConfig struct {
// Brokers is the Kafka bootstrap broker list. Required.
Brokers []string
// ImpressionTopic defaults to DefaultImpressionTopic if empty.
ImpressionTopic string
// ConversionTopic defaults to DefaultConversionTopic if empty.
ConversionTopic string
// SaramaConfig, if nil, defaults to a config with
// Producer.Return.Successes=true (required for SyncProducer) and
// Producer.RequiredAcks=WaitForLocal. A caller-supplied SaramaConfig
// must set Producer.Return.Successes itself — sarama.NewSyncProducer
// errors explicitly if it doesn't, rather than this constructor
// silently patching caller-owned config.
SaramaConfig *sarama.Config
// Logger receives optional diagnostic messages. A nil Logger disables
// logging.
Logger Logger
}
KafkaAnalyticsPublisherConfig configures an AnalyticsPublisher constructed by NewKafkaAnalyticsPublisher.
type KafkaConsumerConfig ¶
type KafkaConsumerConfig struct {
// Brokers is the Kafka bootstrap broker list. Required.
Brokers []string
// GroupID is the consumer group ID. Required.
GroupID string
// Topics are the topics to subscribe to. Required, non-empty.
Topics []string
// SaramaConfig, if nil, defaults to a config matching the reference
// implementation's own defaults (sarama.V3_0_0_0, round-robin
// rebalancing, OffsetOldest, Consumer.Return.Errors=true).
SaramaConfig *sarama.Config
// Logger receives optional diagnostic messages. A nil Logger disables
// logging.
Logger Logger
}
KafkaConsumerConfig configures an EventConsumer constructed by NewKafkaEventConsumer.
type LocaleResolver ¶
type LocaleResolver interface {
ResolveLocale(ctx context.Context, userID string) (string, error)
ResolveLocaleForAnonymous(ctx context.Context, anonymousID string) (string, error)
GetDefaultLocale() string
}
LocaleResolver determines which locale to render a notification in.
func NewPreferencesLocaleResolver ¶
func NewPreferencesLocaleResolver(store PreferencesStore, fallbackLocale string) LocaleResolver
NewPreferencesLocaleResolver constructs a LocaleResolver backed by store.
Parameters:
- store: PreferencesStore
- fallbackLocale: string — used when a user has no stored locale preference, or store lookup fails; defaults to "en" if empty
func NewStaticLocaleResolver ¶
func NewStaticLocaleResolver(locale string) LocaleResolver
NewStaticLocaleResolver returns a LocaleResolver that always resolves to locale, for tests or single-language applications.
type LocalizationStore ¶
type LocalizationStore interface {
// GetLocalizedTemplate returns eventType's template in locale, falling
// back to the LocalizedTemplate's own DefaultLocale if locale isn't
// registered for eventType.
//
// Returns:
// - error: wraps ErrTemplateNotFound if eventType has no
// LocalizedTemplate registered at all
GetLocalizedTemplate(eventType EventType, locale string) (MessageTemplate, error)
RegisterLocalizedTemplate(eventType EventType, locale string, template MessageTemplate) error
// GetSupportedLocales returns every locale registered for eventType,
// or an empty (non-nil) slice if none are.
GetSupportedLocales(eventType EventType) []string
}
LocalizationStore holds per-locale MessageTemplate variants.
func NewInMemoryLocalizationStore ¶
func NewInMemoryLocalizationStore() LocalizationStore
NewInMemoryLocalizationStore constructs an in-memory LocalizationStore.
type LocalizedTemplate ¶
type LocalizedTemplate struct {
DefaultLocale string
Templates map[string]MessageTemplate // locale -> template
}
LocalizedTemplate holds per-locale MessageTemplate variants for one EventType, plus which locale to fall back to when a requested locale isn't registered.
type Logger ¶
type Logger interface {
Debug(msg string, args ...any)
Info(msg string, args ...any)
Warn(msg string, args ...any)
Error(msg string, args ...any)
}
Logger is the minimal logging interface grnoti accepts for optional diagnostic logging (backend connectivity failures, dispatch retries, circuit-breaker state transitions, shutdown). Its four methods match *slog.Logger's own signatures exactly, so *slog.Logger satisfies it structurally — grnoti itself does not import grlog or log/slog, so plugging in a logger is entirely opt-in and adds no dependency for consumers who don't want one.
A nil Logger passed to any constructor is replaced with NopLogger() — logging is always optional, never required for grnoti to function.
Example, using grlog via its log/slog adapter (the recommended bridge — grlog itself needs no code changes for this):
import ( "log/slog" "github.com/gourdian25/grlog" ) logger := slog.New(grlog.NewSlogHandler(grlog.NewDefaultLogger())) deps.Logger = logger svc, err := grnoti.NewNotificationService(deps)
func NopLogger ¶
func NopLogger() Logger
NopLogger returns a Logger that discards every message. It is the default used whenever no Logger is configured.
Returns:
- Logger: a non-nil, no-op implementation safe to call from any goroutine
type Message ¶
type Message struct {
Title string
Body string
Data map[string]string
ImageURL string
Priority Priority
TTL time.Duration
CollapseKey string
ChannelID string
Badge *int
Sound string
Actions []NotificationAction
DeepLink string
Category NotificationCategory
}
Message is the fully-rendered, platform-agnostic notification content a PushDispatcher sends.
type MessageTemplate ¶
type MessageTemplate struct {
TitleTemplate string
BodyTemplate string
DefaultData map[string]string
DefaultTTL time.Duration
CollapseKey string
ChannelID string
Sound string
Actions []NotificationAction
DeepLink string
Category NotificationCategory
}
MessageTemplate is what callers register with a TemplateEngine to define how an EventType renders into a Message.
type Metrics ¶
type Metrics interface {
IncNotificationsProcessed()
IncNotificationsSent(eventType EventType, platform Platform, count int)
IncNotificationsFailed(eventType EventType, platform Platform, count int)
IncInvalidTokens(count int)
IncEventsSkipped(reason string)
ObserveDispatchLatency(eventType EventType, platform Platform, duration time.Duration)
ObserveProcessingLatency(duration time.Duration)
}
Metrics receives counters/observations from dispatch and processing. Unlike the reference implementation, the by-type/by-platform variants take both labels together in one call rather than three separate methods that each only populate one label dimension — see docs/plan/grnoti-plan.md §2 item 10 for why the split version triple-counted.
type MongoDLQHandlerConfig ¶
type MongoDLQHandlerConfig struct {
URI string
Database string
CollectionName string // defaults to DefaultDLQCollection
MaxRetries int // defaults to 3
RetryDelay time.Duration // passed through as-is; 0 means immediately retry-eligible
MaxRetryDelay time.Duration // passed through as-is to FullJitterBackoff (0 there means its own internal default ceiling)
Logger Logger
}
MongoDLQHandlerConfig configures a DLQHandler constructed by NewMongoDLQHandler.
type MongoTokenStoreConfig ¶
type MongoTokenStoreConfig struct {
// URI is a standard MongoDB connection string. Required.
URI string
// Database is the database name. Required.
Database string
// CollectionName defaults to DefaultTokenCollection if empty.
CollectionName string
// Logger receives optional diagnostic messages. A nil Logger disables
// logging.
Logger Logger
}
MongoTokenStoreConfig configures a TokenStore constructed by NewMongoTokenStore. Following grcache's pattern (not gourdiantoken's) — see docs/plan/grnoti-plan.md §1.5 — the store owns and connects its own *mongo.Client from URI, rather than taking an already-connected *mongo.Database from the caller.
type NotificationAction ¶
type NotificationAction struct {
ID string `json:"id"`
Title string `json:"title"`
Icon string `json:"icon,omitempty"`
URL string `json:"url,omitempty"`
}
NotificationAction is one rich-push action button.
type NotificationCategory ¶
type NotificationCategory string
NotificationCategory classifies a Message for filtering/preference purposes (e.g. a user opting out of marketing but not transactional notifications).
const ( // CategoryTransactional covers account/security/order-status content a // user generally cannot opt out of independent of other categories. CategoryTransactional NotificationCategory = "transactional" // CategoryMarketing covers promotional content — the category users // most commonly disable via NotificationPreferences.EventTypeSettings. CategoryMarketing NotificationCategory = "marketing" // CategorySocial covers social-graph activity (likes, follows, mentions). CategorySocial NotificationCategory = "social" // CategoryAlert covers system/operational alerts (see the // EventTypeSystemAlert default template in templateengine.go). CategoryAlert NotificationCategory = "alert" )
type NotificationFailedPayload ¶
type NotificationFailedPayload struct {
EventID string
UserID string
AnonymousID string
EventType EventType
FailureCount int
Reason string
Timestamp time.Time
}
NotificationFailedPayload is published on TopicNotificationFailed.
type NotificationPreferences ¶
type NotificationPreferences struct {
UserID string
GlobalEnabled bool
QuietHoursEnabled bool
QuietHoursStart string // "HH:MM", in Timezone
QuietHoursEnd string // "HH:MM", in Timezone
Timezone string // IANA tz name, e.g. "America/New_York"
Locale string
EventTypeSettings map[EventType]bool
CreatedAt time.Time
UpdatedAt time.Time
}
NotificationPreferences holds one user's notification settings.
func (*NotificationPreferences) IsEventTypeEnabled ¶
func (p *NotificationPreferences) IsEventTypeEnabled(eventType EventType) bool
IsEventTypeEnabled evaluates whether eventType should be sent under p: false if globally disabled, else the explicit per-type setting if one exists, else true (an unconfigured event type defaults to enabled, not disabled). Shared by every PreferencesStore implementation's own IsEventTypeEnabled method so the "unconfigured defaults to enabled" rule is defined exactly once.
type NotificationSentPayload ¶
type NotificationSentPayload struct {
EventID string
UserID string
AnonymousID string
EventType EventType
SuccessCount int
FailureCount int
Timestamp time.Time
}
NotificationSentPayload is published on TopicNotificationSent.
type NotificationService ¶
type NotificationService interface {
ProcessEvent(ctx context.Context, event Event) (ProcessingResult, error)
// Submit is the ingestion-bridge entrypoint (docs/plan/grnoti-plan.md
// §3.1): pass it directly as an EventConsumer's handler —
// consumer.Start(ctx, service.Submit) — to wire Kafka ingestion
// straight into this service. When the service was constructed with
// ServiceConfig.EnableBackpressure, Submit enqueues onto the
// service's own bounded WorkerPool (non-blocking, ErrWorkerPoolFull
// on a full queue) instead of processing on the caller's goroutine;
// otherwise it's equivalent to ProcessEvent with the ProcessingResult
// discarded. Its signature matches WorkerPool's own Handler shape
// exactly, which is why it has no ProcessingResult return — use
// ProcessEvent directly when the caller needs that.
Submit(ctx context.Context, event Event) error
// Close stops any background workers (see WorkerPool) this service
// owns and releases resources. Idempotent. The reference
// implementation's NotificationService had no Close at all — see
// docs/plan/grnoti-plan.md §3.1/§3.6 for why this service now owns a
// WorkerPool that needs one.
Close() error
}
NotificationService is the top-level orchestrator: validates an Event, checks preferences and idempotency, renders it, resolves recipients, dispatches, and records the outcome (metrics, invalid-token cleanup, DLQ on exhausted failure).
func NewNotificationService ¶
func NewNotificationService(deps ServiceDeps) (NotificationService, error)
NewNotificationService constructs a NotificationService.
Parameters:
- deps: ServiceDeps — TokenStore, Dispatcher, Templates, Idempotency are required; everything else is optional and nil-safe
Returns:
- NotificationService: ready to use immediately; if deps.Config.EnableBackpressure is set, its internal WorkerPool is already started
- error: non-nil if a required dependency is missing, or if building the internal WorkerPool fails
type NotificationTarget ¶
type NotificationTarget interface {
IsTopicBased() bool
GetTopicName() string
GetTokens() []DeviceToken
}
NotificationTarget is where a resolved Event should be sent — either a fixed set of device tokens, or an FCM topic.
type PayloadValidator ¶
type PayloadValidator interface {
// ValidateSize returns ErrFCMPayloadTooLarge if msg's estimated
// serialized size exceeds FCMMaxPayloadSize.
ValidateSize(msg Message) error
// EstimateSize returns msg's estimated serialized size in bytes.
EstimateSize(msg Message) int
}
PayloadValidator checks a Message against FCM's payload size limit before attempting to send it.
func NewFCMPayloadValidator ¶
func NewFCMPayloadValidator() PayloadValidator
NewFCMPayloadValidator returns a stateless PayloadValidator that estimates a Message's serialized FCM payload size.
type Platform ¶
type Platform string
Platform is the device platform a DeviceToken/DispatchResult applies to.
const ( // PlatformAndroid is also the fallback platform resolveTokensForEvent // (topicrouter.go) assigns to direct tokens supplied on an Event — bare // strings with no platform information of their own. PlatformAndroid Platform = "android" PlatformIOS Platform = "ios" PlatformWeb Platform = "web" )
type PostgresConfig ¶
type PostgresConfig struct {
// DSN is a standard libpq/pgx connection string. Exactly one of DSN
// or Pool must be set.
DSN string
// Pool, if set, is used directly instead of dialing a new pool from
// DSN — lets multiple grnoti Postgres stores (or the rest of your
// backend) share one pgxpool.Pool instead of each store opening its
// own. grnoti never closes a Pool it did not create itself: every
// store's Close() only closes the pool when it was dialed from DSN.
// Exactly one of DSN or Pool must be set. See docs/postgres.md for
// the recommended shared-pool pattern.
Pool *pgxpool.Pool
// MaxConns caps the pgxpool connection pool size. 0 means use pgxpool's
// own default. Ignored when Pool is set — tune the pool yourself
// before passing it in.
MaxConns int32
// MinConns keeps at least this many connections open. 0 means use
// pgxpool's own default. Ignored when Pool is set.
MinConns int32
// MaxConnLifetime bounds how long a pooled connection may be reused
// before being recycled. 0 means pgxpool's own default (unlimited).
// Ignored when Pool is set.
MaxConnLifetime time.Duration
// ConnectTimeout bounds dialing and the initial Ping when connecting
// from DSN. 0 means 10 seconds. Ignored when Pool is set (the Ping
// against an already-established Pool uses the 10-second default
// unconditionally).
ConnectTimeout time.Duration
// Logger receives optional diagnostic messages. A nil Logger disables
// logging.
Logger Logger
}
PostgresConfig is the common connection configuration shared by every Postgres-backed store.
type PostgresDLQHandlerConfig ¶
type PostgresDLQHandlerConfig struct {
PostgresConfig
MaxRetries int // defaults to 3
RetryDelay time.Duration // 0 is a valid, deliberate "immediately retry-eligible" choice — not defaulted
MaxRetryDelay time.Duration // passed through to FullJitterBackoff as-is
}
PostgresDLQHandlerConfig configures a DLQHandler constructed by NewPostgresDLQHandler — the primary DLQ backend (see docs/plan/grnoti-plan.md §1.3, §6).
type PreferencesFilter ¶
type PreferencesFilter interface {
// ShouldSendNotification evaluates event against its target user's
// preferences.
//
// Returns:
// - bool: false means "do not send"
// - string: a short machine-readable reason when bool is false (e.g.
// "quiet_hours", "global_disabled", "event_type_disabled"),
// surfaced as ProcessingResult.SkipReason; empty when bool is true
// - error: a genuine operational failure (e.g. PreferencesStore
// unreachable) — implementations should fail open (allow the send)
// rather than silently drop a notification when preferences can't
// be evaluated, and still return the error so the caller can log it
ShouldSendNotification(ctx context.Context, event Event) (bool, string, error)
}
PreferencesFilter decides whether an Event should be sent at all, given the target user's preferences (global toggle, quiet hours, per-type opt-out).
func NewPreferencesFilter ¶
func NewPreferencesFilter(store PreferencesStore, logger Logger) PreferencesFilter
NewPreferencesFilter constructs the default PreferencesFilter, evaluating global enable/disable, per-event-type opt-out, and quiet hours against store.
type PreferencesStore ¶
type PreferencesStore interface {
// GetPreferences returns userID's preferences.
//
// Returns:
// - error: satisfies errors.Is(err, ErrPreferencesNotFound) if none
// exist yet — callers generally treat this as "use defaults," not a
// hard failure; check via errors.Is, not direct equality, since an
// implementation may wrap it with additional context
GetPreferences(ctx context.Context, userID string) (*NotificationPreferences, error)
// SavePreferences upserts prefs. prefs.UserID must be non-empty (see
// ErrPreferencesUserIDRequired).
SavePreferences(ctx context.Context, prefs *NotificationPreferences) error
// IsEventTypeEnabled reports whether userID should receive
// notifications of eventType, applying defaults ("enabled") when the
// user has no preferences record or no explicit setting for
// eventType — an unconfigured user is opted in, not opted out.
IsEventTypeEnabled(ctx context.Context, userID string, eventType EventType) (bool, error)
// Close releases any underlying connections. Idempotent.
Close() error
}
PreferencesStore persists per-user notification preferences (global on/off, quiet hours, per-event-type opt-out, locale).
func NewCachedPreferencesStore ¶
func NewCachedPreferencesStore(store PreferencesStore, cache grcache.Cache, ttl time.Duration, logger Logger) PreferencesStore
NewCachedPreferencesStore wraps store with a grcache.Cache-backed read-through cache.
Parameters:
- store: PreferencesStore — the durable source of truth
- cache: grcache.Cache — caller-owned; not closed by this store's Close
- ttl: time.Duration — cache entry TTL; 0 means no expiry (rely entirely on tag invalidation to keep entries fresh)
- logger: Logger — may be nil
func NewMemoryPreferencesStore ¶
func NewMemoryPreferencesStore() PreferencesStore
NewMemoryPreferencesStore constructs an in-memory PreferencesStore.
func NewPostgresPreferencesStore ¶
func NewPostgresPreferencesStore(cfg PostgresConfig) (PreferencesStore, error)
NewPostgresPreferencesStore connects per cfg — the source-of-truth PreferencesStore backend (see docs/plan/grnoti-plan.md §6; pair with NewCachedPreferencesStore for the Redis read-through cache in front of it).
type Priority ¶
type Priority string
Priority is a notification's delivery priority.
const ( // PriorityHigh requests immediate delivery (FCM: high priority — wakes // a doze-mode Android device, bypasses APNs' throttling). PriorityHigh Priority = "high" // PriorityNormal is standard, battery-friendly delivery. PriorityNormal Priority = "normal" // PriorityLow is deliberately deprioritized (e.g. marketing content). // dispatcher.fcm.go's Android/APNS config only branch on == PriorityHigh // (so Low behaves like Normal there), but Webpush maps it to a distinct // "Urgency: low" header — see buildWebpushConfig. PriorityLow Priority = "low" )
type ProcessingResult ¶
type ProcessingResult struct {
EventID string
UserID string
TokenCount int
DispatchResult DispatchResult
ProcessedAt time.Time
Duration time.Duration
Skipped bool
SkipReason string
}
ProcessingResult summarizes the outcome of NotificationService.ProcessEvent.
type PushDispatcher ¶
type PushDispatcher interface {
// Send delivers msg to every token in tokens, batching/fanning out by
// platform internally. A partial failure (some tokens succeed, others
// don't) is reported via the returned DispatchResult, not a non-nil
// error — Send returns a non-nil error only when it could not attempt
// delivery at all (e.g. msg fails payload-size validation).
Send(ctx context.Context, tokens []DeviceToken, msg Message) (DispatchResult, error)
// SendToToken sends msg to a single token, with no batching.
SendToToken(ctx context.Context, token DeviceToken, msg Message) error
// SendToTopic sends msg to every device subscribed to topic via FCM's
// own topic-messaging feature. FCM does not report per-recipient
// results for topic sends, so success/failure here means "the FCM API
// call itself succeeded/failed," not delivery to any specific device.
SendToTopic(ctx context.Context, topic string, msg Message) error
}
PushDispatcher sends rendered Messages to devices or topics via FCM.
func NewFCMDispatcher ¶
func NewFCMDispatcher(deps FCMDispatcherDeps) (PushDispatcher, error)
NewFCMDispatcher constructs an FCM-backed PushDispatcher.
Parameters:
- deps: FCMDispatcherDeps — deps.Client is required
Returns:
- PushDispatcher
- error: ErrFCMClientNil if deps.Client is nil; non-nil if deps.Config.EnableRetry is true and MaxRetryAttempts <= 0
type RateLimiter ¶
type RateLimiter interface {
// Allow reports whether a request may proceed right now, without
// blocking. Consumes a token if true.
Allow(ctx context.Context) (bool, error)
// Wait blocks until a token is available or ctx is done.
Wait(ctx context.Context) error
// GetStats returns a point-in-time snapshot of this limiter's counters.
GetStats(ctx context.Context) (RateLimiterStats, error)
}
RateLimiter bounds outbound FCM request rate. See ratelimiter.go (local, per-process) and ratelimiter.redis.go (distributed) for the two implementations — deliberately different backends behind one interface, see docs/plan/grnoti-plan.md §1.1.
The interface deliberately has no Close(): localRateLimiter owns no resource, so requiring it would force a no-op on every implementation. redisRateLimiter does own a *redis.Client and exposes a Close() error method on its concrete type — callers using the Redis-backed variant type-assert to it (or to an io.Closer) when they need to shut it down, the same pattern already used for UpdateLimit.
func NewLocalRateLimiter ¶
func NewLocalRateLimiter(requestsPerSecond, burstSize int) (RateLimiter, error)
NewLocalRateLimiter constructs a per-process RateLimiter.
Parameters:
- requestsPerSecond: int — must be > 0
- burstSize: int — must be >= requestsPerSecond
Returns:
- RateLimiter
- error: non-nil if either constraint is violated
func NewRedisRateLimiter ¶
func NewRedisRateLimiter(cfg RedisRateLimiterConfig) (RateLimiter, error)
NewRedisRateLimiter builds its own *redis.Client from cfg and validates connectivity with a Ping before returning, mirroring grcache/redis's constructor-time validation.
Parameters:
- cfg: RedisRateLimiterConfig — Addr, RequestsPerSecond, and BurstSize are required; other fields default (see field docs)
Returns:
- RateLimiter: ready to use, shared across every process using the same cfg.Addr/cfg.Key pair
- error: non-nil if a required field is invalid or the connection fails
type RateLimiterStats ¶
type RateLimiterStats struct {
RequestsPerSecond int
BurstSize int
AllowedCount int64
BlockedCount int64
WaitCount int64
LastAllowedAt time.Time
}
RateLimiterStats is a point-in-time snapshot of a RateLimiter's counters.
type RedisRateLimiterConfig ¶
type RedisRateLimiterConfig struct {
// Addr is the Redis server address, e.g. "localhost:6379". Required.
Addr string
// Password authenticates with the server. Empty means no auth.
Password string
// DB selects the Redis logical database.
DB int
// PoolSize is the maximum number of connections in the pool. Defaults to 100.
PoolSize int
// DialTimeout bounds how long connecting to Redis may take. Defaults to 5s.
DialTimeout time.Duration
// ReadTimeout bounds how long a read may take. Defaults to 3s.
ReadTimeout time.Duration
// WriteTimeout bounds how long a write may take. Defaults to 3s.
WriteTimeout time.Duration
// RequestsPerSecond is the bucket's steady-state refill rate, shared
// across every process using the same Key. Required, must be > 0.
RequestsPerSecond int
// BurstSize is the bucket's capacity. Required, must be >= RequestsPerSecond.
BurstSize int
// Key identifies the shared bucket. All processes that should share
// one distributed quota must use the same Key. Defaults to
// "grnoti:ratelimit:default".
Key string
// Logger receives optional diagnostic messages. A nil Logger disables logging.
Logger Logger
}
RedisRateLimiterConfig configures a redisRateLimiter constructed by NewRedisRateLimiter. Zero-valued connection fields fall back to the same defaults as grcache/redis's RedisConfig; RequestsPerSecond and BurstSize have no sensible zero value and must be set explicitly.
type RetryStrategy ¶
type RetryStrategy interface {
// ShouldRetry reports whether attempt (0-indexed) should be retried
// given err.
ShouldRetry(attempt int, err error) bool
// GetDelay returns how long to wait before attempt (0-indexed).
GetDelay(attempt int) time.Duration
}
RetryStrategy decides whether and how long to wait before retrying a failed FCM send.
func NewFullJitterRetry ¶
func NewFullJitterRetry(maxAttempts int, baseDelay, maxDelay time.Duration) RetryStrategy
NewFullJitterRetry constructs a RetryStrategy for FCM dispatch retries, backed by FullJitterBackoff. Replaces the reference implementation's un-jittered base*2^attempt strategy (see docs/plan/grnoti-plan.md §1.2) — a synchronized retry storm across replicas after an FCM outage is exactly what jitter exists to avoid.
Parameters:
- maxAttempts: int — total attempts, including the first; ShouldRetry returns false once this many attempts have been made
- baseDelay, maxDelay: time.Duration — passed through to FullJitterBackoff
func NewNoopRetryStrategy ¶
func NewNoopRetryStrategy() RetryStrategy
NewNoopRetryStrategy returns a RetryStrategy that never retries, for tests and for dispatchers explicitly configured without retry.
type ServiceConfig ¶
type ServiceConfig struct {
// IdempotencyTTL is how long a processed event's ID is remembered by
// IdempotencyStore before it could be reprocessed if redelivered again.
IdempotencyTTL time.Duration
// MaxTokensPerBatch is the batch size dispatchToTokens splits a token
// list into when EnforceBatching is set. Distinct from (and composes
// with) dispatcher.fcm.go's own fixed FCMMaxBatchSize internal batching.
MaxTokensPerBatch int
// EnableMetrics gates whether processEvent/recordMetrics call the
// configured Metrics collector at all.
EnableMetrics bool
// SkipInvalidEvents, if set, makes processEvent report a failed
// Event.Validate() as a skip (ProcessingResult.Skipped, no error
// returned) instead of propagating the validation error to the caller.
SkipInvalidEvents bool
// EnableTokenDeduplication makes dispatchToTokens deduplicate the
// resolved token list (via batchSplitter.Deduplicate) before sending,
// in case TokenStore ever returns the same token twice.
EnableTokenDeduplication bool
// EnforceBatching makes dispatchToTokens pre-split a resolved token
// list into MaxTokensPerBatch-sized calls to the dispatcher, merging
// the per-batch DispatchResults — off by default, in which case the
// entire token list is passed to the dispatcher in one call.
EnforceBatching bool
// EnablePreferencesFilter gates whether processEvent consults
// ServiceDeps.PreferencesFilter for authenticated events at all; a nil
// PreferencesFilter always skips this check regardless of the flag.
EnablePreferencesFilter bool
// EnableTopicRouting gates whether processEvent consults
// ServiceDeps.TopicRouter to resolve a topic/token target instead of
// always resolving tokens directly; a nil TopicRouter always skips
// this regardless of the flag.
EnableTopicRouting bool
// EnableRichPush, EnableLocalization, EnableABTesting are
// composition-time flags, not read anywhere in
// notificationService.processEvent (service.go) directly — rich-push
// fields live on Message and are populated by whichever
// TemplateEngine ServiceDeps.Templates is; localization is a
// TemplateEngine decorator (localizedTemplateEngine, see
// localization.go) the caller wraps ServiceDeps.Templates in; A/B
// assignment (ExperimentEngine.AssignVariant) happens before an Event
// is even constructed, to decide what goes in its Payload. These
// three flags exist for a caller's own bookkeeping/documentation of
// which optional pieces a given ServiceDeps wiring includes, not as
// live branches in the pipeline.
EnableRichPush bool
EnableLocalization bool
EnableBackpressure bool
EnableABTesting bool
// EnableDLQ gates whether ProcessEvent publishes exhausted-retry
// dispatch failures to the configured DLQHandler. The reference
// implementation had no such wiring at all (see
// docs/plan/grnoti-plan.md §3.6) — this flag exists so the fix is
// opt-in-by-default-on rather than silently always-on, matching how
// every other integration point in this config is a flag.
EnableDLQ bool
// EnableEventBus gates whether ProcessEvent/ExperimentEngine publish
// lifecycle events via the configured grevents.Bus (see
// docs/plan/grnoti-plan.md §1.2).
EnableEventBus bool
}
ServiceConfig configures a NotificationService's optional behaviors. All Enable* flags default to false ("disabled by default for predictability and backward compatibility" — matching the reference implementation's own convention) except where DefaultServiceConfig says otherwise.
func DefaultServiceConfig ¶
func DefaultServiceConfig() ServiceConfig
DefaultServiceConfig returns a ServiceConfig with sane production defaults: metrics on, DLQ publishing on, everything else opt-in.
type ServiceDeps ¶
type ServiceDeps struct {
// TokenStore, Dispatcher, Templates, Idempotency are required.
TokenStore TokenStore
Dispatcher PushDispatcher
Templates TemplateEngine
Idempotency IdempotencyStore
// PreferencesFilter, if set and Config.EnablePreferencesFilter, gates
// authenticated-user dispatch on ShouldSendNotification.
PreferencesFilter PreferencesFilter
// TopicRouter, if set and Config.EnableTopicRouting, resolves each
// Event's NotificationTarget instead of the default direct-tokens
// resolution (resolveTokensForEvent in topicrouter.go).
TopicRouter TopicRouter
// DLQHandler, if set and Config.EnableDLQ, receives events whose
// dispatch has unresolved failures after the dispatcher's own retry
// is exhausted — the reference implementation built a whole DLQ
// subsystem nothing ever called (docs/plan/grnoti-plan.md §3.6); this
// is that missing call.
DLQHandler DLQHandler
// EventBus, if set and Config.EnableEventBus, receives
// TopicNotificationSent/TopicNotificationFailed lifecycle events (see
// events.go, §1.2).
EventBus grevents.Bus
// Metrics is optional. NotificationService only calls the
// per-event/per-platform counters (IncNotificationsProcessed/Sent/
// Failed, ObserveDispatchLatency/ProcessingLatency, IncEventsSkipped)
// — it deliberately does NOT call IncInvalidTokens, since
// dispatcher.fcm.go already calls that itself when the same Metrics
// instance is wired into FCMDispatcherDeps.Metrics; calling it again
// here would double-count.
Metrics Metrics
// Config toggles pipeline behavior. See DefaultServiceConfig.
Config ServiceConfig
// WorkerPoolConfig is used only when Config.EnableBackpressure is
// true, to build the service's own internal *WorkerPool (see Submit).
WorkerPoolConfig WorkerPoolConfig
Logger Logger
}
ServiceDeps configures a NotificationService constructed by NewNotificationService.
type TemplateEngine ¶
type TemplateEngine interface {
// BuildMessage renders event into a Message using the MessageTemplate
// registered for event.Type, falling back to the template registered
// under EventTypeCustom if event.Type has none of its own.
//
// Returns:
// - error: wraps ErrTemplateNotFound if neither event.Type nor
// EventTypeCustom has a registered template
BuildMessage(event Event) (Message, error)
RegisterTemplate(eventType EventType, template MessageTemplate) error
}
TemplateEngine renders an Event into a Message using registered MessageTemplates.
func NewLocalizedTemplateEngine ¶
func NewLocalizedTemplateEngine(baseEngine TemplateEngine, localeStore LocalizationStore, localeResolver LocaleResolver) TemplateEngine
NewLocalizedTemplateEngine wraps baseEngine with locale-aware rendering, itself implementing TemplateEngine so it's a drop-in replacement.
func NewTemplateEngine ¶
func NewTemplateEngine() TemplateEngine
NewTemplateEngine constructs a TemplateEngine pre-seeded with a small set of generic default templates (see registerDefaults) — deliberately scheme-free and deep-link-free, unlike the reference implementation, which hardcoded a "skipp://" deep-link scheme into 8 of 9 default templates (see docs/plan/grnoti-plan.md §2 item 8). Consumers register their own application-specific templates (including their own deep-link scheme) via RegisterTemplate.
type TokenStore ¶
type TokenStore interface {
// GetActiveTokens returns every active DeviceToken registered for
// userID.
GetActiveTokens(ctx context.Context, userID string) ([]DeviceToken, error)
// GetActiveTokensBatch is the multi-user form of GetActiveTokens,
// returning a map keyed by userID. A userID with no active tokens is
// simply absent from the result map, not an error.
GetActiveTokensBatch(ctx context.Context, userIDs []string) (map[string][]DeviceToken, error)
// GetActiveTokensByAnonymousID is GetActiveTokens for an anonymous
// (pre-authentication) visitor.
GetActiveTokensByAnonymousID(ctx context.Context, anonymousID string) ([]DeviceToken, error)
// MarkInvalid deactivates token (e.g. after FCM reports it as
// unregistered). Marking an already-inactive or nonexistent token is
// not an error.
MarkInvalid(ctx context.Context, token string) error
// SaveToken upserts token: creates it if new, refreshes its metadata
// and reactivates it (IsActive=true) if it already exists.
SaveToken(ctx context.Context, token DeviceToken) error
// DeleteToken permanently removes token. Deleting a nonexistent token
// is not an error.
DeleteToken(ctx context.Context, token string) error
// Close releases any underlying connections. Idempotent.
Close() error
}
TokenStore manages device-token registration and lookup. Every backend (mongo.go, postgres.go, memory.go) implements this identically so application code depends on the interface alone.
func NewMemoryTokenStore ¶
func NewMemoryTokenStore() TokenStore
NewMemoryTokenStore constructs an in-memory TokenStore.
func NewMongoTokenStore ¶
func NewMongoTokenStore(cfg MongoTokenStoreConfig) (TokenStore, error)
NewMongoTokenStore connects to MongoDB per cfg, ensures indexes, and validates connectivity via Ping before returning.
func NewPostgresTokenStore ¶
func NewPostgresTokenStore(cfg PostgresConfig) (TokenStore, error)
NewPostgresTokenStore connects per cfg (pgx + sqlc-generated queries) — the alternative to the primary MongoDB backend (tokenstore.mongo.go); see docs/plan/grnoti-plan.md §6.
type TokenTarget ¶
type TokenTarget struct {
Tokens []DeviceToken
}
TokenTarget is a NotificationTarget resolving to a fixed set of device tokens.
func (TokenTarget) GetTokens ¶
func (t TokenTarget) GetTokens() []DeviceToken
func (TokenTarget) GetTopicName ¶
func (t TokenTarget) GetTopicName() string
func (TokenTarget) IsTopicBased ¶
func (t TokenTarget) IsTopicBased() bool
type TopicRouter ¶
type TopicRouter interface {
ResolveTarget(ctx context.Context, event Event) (NotificationTarget, error)
}
TopicRouter resolves an Event to a NotificationTarget.
func NewEventTypeTopicRouter ¶
func NewEventTypeTopicRouter(topicMappings map[EventType]string, tokenStore TokenStore, logger Logger) TopicRouter
NewEventTypeTopicRouter constructs the primary TopicRouter.
Parameters:
- topicMappings: map[EventType]string — may be nil
- tokenStore: TokenStore — used as the fallback when neither a payload override nor a type mapping applies
- logger: Logger — may be nil
func NewStaticTopicRouter ¶
func NewStaticTopicRouter(topic string) TopicRouter
NewStaticTopicRouter returns a TopicRouter that always resolves to the same fixed topic, for tests or single-topic applications.
func NewTokenOnlyRouter ¶
func NewTokenOnlyRouter(tokenStore TokenStore) TopicRouter
NewTokenOnlyRouter returns a TopicRouter that always routes via TokenStore, disabling topic-based routing entirely.
type TopicTarget ¶
type TopicTarget struct {
Topic string
}
TopicTarget is a NotificationTarget resolving to an FCM topic.
func (TopicTarget) GetTokens ¶
func (t TopicTarget) GetTokens() []DeviceToken
func (TopicTarget) GetTopicName ¶
func (t TopicTarget) GetTopicName() string
func (TopicTarget) IsTopicBased ¶
func (t TopicTarget) IsTopicBased() bool
type WorkerPool ¶
type WorkerPool struct {
// contains filtered or unexported fields
}
WorkerPool is the bridge between event ingestion (EventConsumer) and event processing (NotificationService.ProcessEvent), decoupling the two via a bounded, non-blocking queue. Unlike the reference implementation (see docs/plan/grnoti-plan.md §3.1), grnoti's EventConsumer and NotificationService are wired through a WorkerPool by default — an ingestion handler that calls ProcessEvent directly, with no queue in between, is exactly the gap this type exists to close.
WorkerPool is deliberately a concrete type, not an interface — it has exactly one implementation and no swappable backend, matching CircuitBreaker/RateLimiter's (local)/RetryStrategy's own treatment.
func NewWorkerPool ¶
func NewWorkerPool(deps WorkerPoolDeps) (*WorkerPool, error)
NewWorkerPool constructs a WorkerPool. Call Start to begin processing and Stop (or Close) to shut down.
Parameters:
- deps: WorkerPoolDeps — deps.Handler must be non-nil; deps.Config.Workers defaults to 10 if <= 0, deps.Config.QueueSize defaults to 1000 if <= 0
Returns:
- *WorkerPool
- error: non-nil if deps.Handler is nil
func (*WorkerPool) GetStats ¶
func (wp *WorkerPool) GetStats() WorkerPoolStats
GetStats returns a point-in-time snapshot of the pool's queue occupancy.
func (*WorkerPool) Start ¶
func (wp *WorkerPool) Start()
Start spawns the pool's worker goroutines. Safe to call at most once.
func (*WorkerPool) Stop ¶
func (wp *WorkerPool) Stop()
Stop closes the queue, waits for every worker to fully drain it and exit, and only then cancels wp.ctx. This ordering is what makes drain deterministic: while the queue is closing but wp.ctx is not yet canceled, ctx.Done() is never ready, so each worker's select has exactly one viable case — read the queue until it reports empty-and- closed. Stop therefore guarantees every event successfully Submitted before Stop was called is delivered to a worker and run through handler before Stop returns. Canceling only after every worker has exited also means handler still runs with a live (non-canceled) ctx for every event drained during shutdown, rather than racing an already-canceled one.
func (*WorkerPool) Submit ¶
func (wp *WorkerPool) Submit(event Event) error
Submit enqueues event without blocking. Returns ErrWorkerPoolFull (wrapped with occupancy detail) if the queue is full — this is backpressure-by-rejection, not backpressure-by-blocking; callers that need to block should retry with their own backoff.
func (*WorkerPool) SubmitAsync ¶
func (wp *WorkerPool) SubmitAsync(event Event) bool
SubmitAsync is Submit with a bool result instead of an error — true if enqueued, false if the queue was full.
type WorkerPoolConfig ¶
type WorkerPoolConfig struct {
// Workers is the number of worker goroutines. Defaults to 10 if <= 0.
Workers int
// QueueSize is the buffered-channel capacity. Defaults to 1000 if <= 0.
QueueSize int
}
WorkerPoolConfig configures a WorkerPool.
type WorkerPoolDeps ¶
type WorkerPoolDeps struct {
Config WorkerPoolConfig
// Handler processes one Event — typically a thin wrapper around
// NotificationService.ProcessEvent. Required.
Handler func(context.Context, Event) error
Logger Logger
// Metrics is optional; if set, IncEventsSkipped("backpressure") is
// called whenever Submit/SubmitAsync rejects an Event for a full queue.
Metrics Metrics
}
WorkerPoolDeps configures a WorkerPool.
Source Files
¶
- batchsplitter.go
- cache.experiment.go
- cache.idempotency.go
- cache.preferences.go
- circuitbreaker.go
- consumer.kafka.go
- dispatcher.fcm.go
- dlq.mongo.go
- dlq.postgres.go
- docs.go
- errors.go
- events.go
- eventtypes.go
- experiment.go
- experimentstore.postgres.go
- interfaces.go
- localization.go
- logger.go
- memory.go
- payloadvalidator.go
- postgres.go
- preferences.postgres.go
- preferencesfilter.go
- producer.kafka.go
- ratelimiter.go
- ratelimiter.redis.go
- retrystrategy.go
- service.go
- templateengine.go
- tokenstore.mongo.go
- tokenstore.postgres.go
- topicrouter.go
- types.go
- version.go
- workerpool.go
Directories
¶
| Path | Synopsis |
|---|---|
|
Command example is a minimal, fully self-contained walkthrough of grnoti: no Docker containers, no Firebase credentials, no network access at all.
|
Command example is a minimal, fully self-contained walkthrough of grnoti: no Docker containers, no Firebase credentials, no network access at all. |
|
internal
|
|