grnoti

package module
v0.1.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Jul 22, 2026 License: MIT Imports: 32 Imported by: 0

README

grnoti

Go Reference Go Version License

Push-notification service library for the gourdian ecosystem (github.com/gourdian25/grnoti): FCM dispatch, idempotent event processing, device-token management, durable dead-letter retry, circuit breaking, distributed rate limiting, deterministic A/B experiment assignment, localization, and topic-based routing, behind a set of storage-agnostic interfaces.

Status: feature-complete per the 14-stage build plan (docs/plan/grnoti-plan.md), pre-tagged-release. golangci-lint run reports 0 issues; test coverage is 95.1% on the root package, enforced by a 95% gate (make coverage-check), verified against real local MongoDB/PostgreSQL/Redis/Kafka instances (see CLAUDE.md for the docker setup).

Table of Contents

Part of the gourdian25 ecosystem

grnoti is one of several small, independent Go libraries meant to be used together:

  • grcache — backend-agnostic caching abstraction; grnoti's NewCacheIdempotencyStore, NewCachedPreferencesStore, and NewCacheBackedExperimentEngine each wrap any grcache.Cache directly, no adapter needed.
  • grevents — an in-process event bus; grnoti optionally publishes notification.sent/ notification.failed/experiment.assigned lifecycle events through it — best-effort, so a nil bus or a publish failure never affects the durable operation it follows.
  • gourdiantoken — JWT access/refresh token issuance, verification, revocation, and rotation.
  • grlog — zero-dependency structured logging.
  • graudit — an append-only, tamper-evident audit log with pluggable storage backends.
  • grpolicy — attribute-based policy evaluation (RBAC/ABAC), independent of any notion of "user" or "role".

Install

go get github.com/gourdian25/grnoti

Dependencies

grnoti is a single flat package with no subpackages of its own (see Why this shape) — every backend lives in the same module, distinguished by a <concern>.<backend>.go file-naming convention. The direct consequence: go get github.com/gourdian25/grnoti pulls in every backend driver, regardless of which ones a given deployment actually uses — the MongoDB driver, pgx/v5 (plus sqlc-generated query code), go-redis, IBM's Sarama (Kafka), and the Firebase Admin SDK. This is a deliberate divergence from sibling repos like grcache/graudit, which use one subpackage per backend specifically to keep unused drivers out of a consumer's dependency graph. grnoti accepts that heavier import in exchange for a simpler package to navigate (one import "github.com/gourdian25/grnoti", not one per backend). If a heavy transitive dependency graph is a concern for your deployment, that tradeoff is worth knowing about up front — it isn't something go mod tidy or build tags can undo here.

Quickstart

templates := grnoti.NewTemplateEngine()
templates.RegisterTemplate("order_shipped", grnoti.MessageTemplate{
    TitleTemplate: "Your order has shipped!",
    BodyTemplate:  "Order #{{.order_id}} is on its way.",
})

tokenStore, err := grnoti.NewMongoTokenStore(grnoti.MongoTokenStoreConfig{URI: mongoURI, Database: "myapp"})
if err != nil {
    log.Fatal(err)
}
dispatcher, err := grnoti.NewFCMDispatcher(grnoti.FCMDispatcherDeps{Client: fcmClient})
if err != nil {
    log.Fatal(err)
}

svc, err := grnoti.NewNotificationService(grnoti.ServiceDeps{
    TokenStore:  tokenStore,
    Dispatcher:  dispatcher,
    Templates:   templates,
    Idempotency: grnoti.NewCacheIdempotencyStore(redisCache), // any grcache.Cache
    Config:      grnoti.DefaultServiceConfig(),
})
if err != nil {
    log.Fatal(err)
}
defer svc.Close()

_, err = svc.ProcessEvent(ctx, grnoti.Event{
    EventID:  "evt-1",
    UserID:   "user-42",
    Type:     "order_shipped",
    Priority: grnoti.PriorityHigh,
    Payload:  map[string]string{"order_id": "1001"},
})

See example/main.go for a complete, runnable, narrated walkthrough — go run ./example, no external services required (it uses the in-memory backends plus a dispatcher that logs to stdout instead of calling FCM). It also documents the exact one-line swap for every real backend constructor.

An intermediate example: DLQ, circuit breaker, rate limiting

The Quickstart above only wires the four required ServiceDeps fields (TokenStore, Dispatcher, Templates, Idempotency). A more realistic production wiring also protects the FCM dispatch path itself and gives failed sends somewhere durable to land. Building on example/main.go's own structure:

// A circuit breaker: after 5 consecutive FCM failures, stop calling FCM
// for 30s and fail fast instead; a closed breaker's failure counter
// resets after 1 minute with no failures.
breaker, err := grnoti.NewCircuitBreaker(5, 30*time.Second, time.Minute)
if err != nil {
    log.Fatal(err)
}

// A local (per-process) rate limiter: at most 50 FCM calls/sec, bursts
// of up to 10. Swap for grnoti.NewRedisRateLimiter(...) to share one
// limit across multiple service instances instead.
limiter, err := grnoti.NewLocalRateLimiter(50, 10)
if err != nil {
    log.Fatal(err)
}

// RateLimiter/CircuitBreaker, once set here, actually gate every
// outbound FCM batch/single-send in dispatcher.fcm.go — not just built
// and left unconnected (docs/architecture.md §3.2).
dispatcher, err := grnoti.NewFCMDispatcher(grnoti.FCMDispatcherDeps{
    Client:         fcmClient,
    Config:         grnoti.DefaultFCMDispatcherConfig(),
    RateLimiter:    limiter,
    CircuitBreaker: breaker,
})
if err != nil {
    log.Fatal(err)
}

// A durable dead-letter queue: any dispatch failure not already
// accounted for by a marked-invalid token is published here instead of
// silently dropped. NewMemoryDLQHandler(maxRetries, retryDelay,
// maxRetryDelay) shown here; swap for NewPostgresDLQHandler/
// NewMongoDLQHandler in production for a restart-durable queue.
dlqHandler := grnoti.NewMemoryDLQHandler(3, time.Minute, 10*time.Minute)

config := grnoti.DefaultServiceConfig() // EnableDLQ is already true by default
svc, err := grnoti.NewNotificationService(grnoti.ServiceDeps{
    TokenStore:  tokenStore,
    Dispatcher:  dispatcher,
    Templates:   templates,
    Idempotency: idempotencyStore,
    DLQHandler:  dlqHandler,
    Config:      config,
})
if err != nil {
    log.Fatal(err)
}
defer svc.Close()

NotificationService never calls DLQHandler.ClaimRetryableEvents itself — it only ever calls PublishToDLQ. Draining the queue is a separate, external retry-worker process's job, polling periodically:

events, err := dlqHandler.ClaimRetryableEvents(ctx, 50) // atomically claims up to 50
// events may be non-nil even when err != nil for some backends (e.g.
// Mongo) — process what was returned regardless; see
// DLQHandler.ClaimRetryableEvents' doc comment (interfaces.go) and
// docs/architecture.md §3.6.
for _, ev := range events {
    // re-attempt delivery, then report the outcome:
    _ = dlqHandler.MarkRetried(ctx, ev.EventID, success, attemptErr)
}

Configuration

NewNotificationService(ServiceDeps) (NotificationService, error)
ServiceDeps field Required? Notes
TokenStore required device-token lookup
Dispatcher required FCM send path
Templates required Event → Message rendering
Idempotency required duplicate-delivery suppression
PreferencesFilter optional gates authenticated dispatch on ShouldSendNotification; only consulted when Config.EnablePreferencesFilter is also set — nil-safe either way
TopicRouter optional resolves each event's NotificationTarget instead of the default direct-token resolution; only consulted when Config.EnableTopicRouting is also set
DLQHandler optional receives unresolved dispatch failures; only consulted when Config.EnableDLQ is also set
EventBus (grevents.Bus) optional receives lifecycle events; only consulted when Config.EnableEventBus is also set — publishing is always best-effort
Metrics optional per-event/per-platform counters and latency observations
WorkerPoolConfig optional only used when Config.EnableBackpressure is set, to build the service's own internal *WorkerPool
Logger optional nil-safe, defaults to a no-op logger
Config (ServiceConfig) optional see table below; zero value is "everything opt-in off" except where DefaultServiceConfig() says otherwise

Missing any of the four required fields makes NewNotificationService return an error immediately — every other field is nil-safe and simply disables the behavior it would have enabled.

ServiceConfig (via DefaultServiceConfig())
Field Default Effect
IdempotencyTTL 24 * time.Hour How long a processed event's ID is remembered by IdempotencyStore before it could be reprocessed if redelivered again
MaxTokensPerBatch 500 (FCM's own per-multicast-call limit) Batch size dispatchToTokens splits a resolved token list into when EnforceBatching is set
EnableMetrics true Gates whether processEvent calls the configured Metrics collector at all
SkipInvalidEvents false A failed Event.Validate() is reported as a skip (ProcessingResult.Skipped) instead of returned as an error
EnableTokenDeduplication false Deduplicates a resolved token list (by Token value) before dispatch
EnforceBatching false Pre-splits a resolved token list into MaxTokensPerBatch-sized dispatcher calls, merging the per-batch results — a coarser, caller-controlled layer that composes with (doesn't replace) the dispatcher's own fixed internal FCM batching
EnablePreferencesFilter false Consults ServiceDeps.PreferencesFilter for authenticated events; no-op if PreferencesFilter is nil regardless of this flag
EnableTopicRouting false Consults ServiceDeps.TopicRouter to resolve a topic/token target instead of always resolving tokens directly; no-op if TopicRouter is nil regardless of this flag
EnableRichPush false Documentation-only. Not read anywhere in processEvent — rich-push fields live on Message and are populated by whichever TemplateEngine is wired in, independent of this flag
EnableLocalization false Documentation-only. Localization is a TemplateEngine decorator (NewLocalizedTemplateEngine) the caller wraps ServiceDeps.Templates in — this flag doesn't gate anything itself
EnableABTesting false Documentation-only. A/B variant assignment (ExperimentEngine.AssignVariant) happens before an Event is even constructed, to decide what goes into its Payload — this flag doesn't gate anything itself
EnableBackpressure false Routes Submit through the service's own internal, bounded *WorkerPool (non-blocking, ErrWorkerPoolFull on a full queue) instead of processing on the caller's goroutine
EnableDLQ true Publishes unresolved dispatch failures to the configured DLQHandler; no-op if DLQHandler is nil regardless of this flag
EnableEventBus false Publishes notification.sent/notification.failed lifecycle events via the configured EventBus; no-op if EventBus is nil regardless of this flag

EnableRichPush, EnableLocalization, and EnableABTesting are real fields on the struct and are safe to set, but as of the current code they are composition-time bookkeeping only — nothing in notificationService.processEvent (service.go) branches on them. What they'd describe is decided entirely by which concrete TemplateEngine/ ExperimentEngine implementation you compose into ServiceDeps before constructing the service, not by anything the pipeline itself can gate on at request time. Don't rely on flipping one of these three to change runtime behavior.

Public API overview

The tables below are constructor-oriented; see docs/architecture.md §2 for the full interface-by-backend matrix (with checkmarks) and the reasoning behind each implementation choice.

TokenStore

Backend Constructor
In-memory NewMemoryTokenStore()
MongoDB NewMongoTokenStore(MongoTokenStoreConfig)
PostgreSQL NewPostgresTokenStore(PostgresConfig)

PreferencesStore / PreferencesFilter

Variant Constructor
In-memory store NewMemoryPreferencesStore()
PostgreSQL store NewPostgresPreferencesStore(PostgresConfig)
Cache-backed read-through (wraps any PreferencesStore) NewCachedPreferencesStore(store, grcache.Cache, ttl, Logger)
Filter (consumes a PreferencesStore, used via ServiceDeps.PreferencesFilter) NewPreferencesFilter(store, Logger)

DLQHandler

Backend Constructor
In-memory NewMemoryDLQHandler(maxRetries, retryDelay, maxRetryDelay)
MongoDB NewMongoDLQHandler(MongoDLQHandlerConfig)
PostgreSQL NewPostgresDLQHandler(PostgresDLQHandlerConfig)

ExperimentStore / ExperimentEngine

Variant Constructor
Store: in-memory NewMemoryExperimentStore()
Store: PostgreSQL NewPostgresExperimentStore(PostgresConfig)
Engine: in-process assignment cache NewDeterministicExperimentEngine(AnalyticsPublisher, grevents.Bus, Logger)
Engine: shared, cache-backed assignment NewCacheBackedExperimentEngine(grcache.Cache, AnalyticsPublisher, grevents.Bus, Logger)

Assignment itself (deterministicPick) is a pure, deterministic function of (userID, experiment.ID, experiment.Variants) in both engines — the cache is purely an optimization for where repeated assignments are remembered, never a correctness dependency.

IdempotencyStore

Backend Constructor
Any grcache.Cache (Redis, in-memory, ...) NewCacheIdempotencyStore(grcache.Cache)

PushDispatcher (FCM)

Backend Constructor
Firebase Cloud Messaging NewFCMDispatcher(FCMDispatcherDeps) — tune retry via FCMDispatcherConfig/DefaultFCMDispatcherConfig() (EnableRetry, MaxRetryAttempts, RetryBaseDelay/RetryMaxDelay)

RateLimiter

Variant Constructor
Local (per-process) NewLocalRateLimiter(requestsPerSecond, burstSize int)
Redis (distributed) NewRedisRateLimiter(RedisRateLimiterConfig)

CircuitBreaker

Variant Constructor
The one implementation (standardCircuitBreaker), returned as the CircuitBreaker interface NewCircuitBreaker(maxFailures, timeout, resetTimeout) or NewCircuitBreakerWithConfig(CircuitBreakerConfig)

Kafka

Role Constructor
EventConsumer NewKafkaEventConsumer(KafkaConsumerConfig)
AnalyticsPublisher NewKafkaAnalyticsPublisher(KafkaAnalyticsPublisherConfig)

Templates, localization, routing, misc.

Capability Constructor
TemplateEngine (scheme-free defaults) NewTemplateEngine()
TemplateEngine (locale-aware decorator) NewLocalizedTemplateEngine(base, LocalizationStore, LocaleResolver)
LocalizationStore NewInMemoryLocalizationStore()
LocaleResolver (reads a user's stored locale) NewPreferencesLocaleResolver(store, fallbackLocale)
LocaleResolver (fixed locale) NewStaticLocaleResolver(locale)
TopicRouter (by event type) NewEventTypeTopicRouter(topicMappings, TokenStore, Logger)
TopicRouter (fixed topic) NewStaticTopicRouter(topic)
TopicRouter (tokens only, no topic routing) NewTokenOnlyRouter(TokenStore)
BatchSplitter NewBatchSplitter()
RetryStrategy (full-jitter backoff) NewFullJitterRetry(maxAttempts, baseDelay, maxDelay)
RetryStrategy (never retries) NewNoopRetryStrategy()
PayloadValidator (FCM payload-size check) NewFCMPayloadValidator()
*WorkerPool (used internally when Config.EnableBackpressure is set) NewWorkerPool(WorkerPoolDeps)

Postgres: sharing one pool across stores

NewPostgresTokenStore, NewPostgresPreferencesStore, NewPostgresExperimentStore, and NewPostgresDLQHandler each take a PostgresConfig. Giving each one its own DSN means each dials its own pool — fine for one store, wasteful for four (MaxConns × 4 connections, not × 1). PostgresConfig.Pool lets you build one *pgxpool.Pool yourself and inject it into every store instead:

pool, err := pgxpool.NewWithConfig(ctx, poolCfg) // your own bootstrap code
if err != nil {
    log.Fatal(err)
}
defer pool.Close() // grnoti never closes a Pool it didn't dial itself

tokenStore, err := grnoti.NewPostgresTokenStore(grnoti.PostgresConfig{Pool: pool})
preferencesStore, err := grnoti.NewPostgresPreferencesStore(grnoti.PostgresConfig{
    Pool: pool, SkipSchemaEnsure: true, // schema already applied by tokenStore above
})

DSN and Pool are mutually exclusive — set exactly one. PostgresConfig.SkipSchemaEnsure skips grnoti's built-in schema application for stores managed by your own migration pipeline instead. See docs/postgres.md for the full pattern, Close() ownership rules, and the concurrency-safety guarantee (schema application is now serialized via a Postgres advisory lock).

Why storage-agnostic interfaces

Every capability — token storage, preferences, dead-letter retry, rate limiting, experiment assignment — is a small interface with real, independently-usable implementations: in-memory (for tests or small deployments), MongoDB, PostgreSQL, or a grcache.Cache-backed adapter (works with any of grcache's own backends, including Redis). Pick whichever combination matches your infrastructure; nothing in NotificationService assumes a specific one. See docs/architecture.md for the full interface/backend matrix and the reasoning behind each design decision.

Why this shape

grnoti was the first repo in the gourdian25 ecosystem to use this shape — a single flat package with pgx/v5 + sqlc-generated Postgres queries and no GORM, ever — rather than adopting it from a sibling repo. grcache, graudit, and gourdiantoken later converged on the same flat-package, GORM-free pattern during their own standardization passes, so grnoti's shape became the template the rest of the ecosystem adopted, not the other way around. See docs.go's "Package shape" section for the full reasoning.

Testing

make test              # go test -count=1 -timeout=5m -cover ./...
make race               # go test -race -timeout 5m ./... — mandatory before any commit
                         # touching experiment.go, workerpool.go, dlq.*.go, or any store
make coverage-summary    # per-function coverage breakdown
make coverage-check      # gates the root package at 95%
make precommit           # fmt + vet + lint + race + coverage-check — run before every commit

Coverage scoping caveat: use go test -cover . (a single dot), not go test -cover ./.... The root package's own coverage is what coverage-check gates on and reports (95.1% at last check), but running against ./... also compiles and instruments internal/postgresdb (sqlc- generated query wrappers with no test file of their own, so it always reports a flat 0%) and the example command package (also untested by design — it's a runnable demo, not library code) — both dilute the printed number without reflecting anything about this package's actual test coverage. make coverage-check already scopes to . for this reason; internal/postgresdb's 0% line in go tool cover -func output can be ignored — see docs/architecture.md §5.

Real local backends, not mocks: every backend (MongoDB, PostgreSQL, Redis, Kafka) is tested against a real local instance via Docker, not a mock — this is a genuine differentiator, and it has caught real bugs mocked tests would have missed (see docs/architecture.md §4 for three concrete examples, including a confirmed data race and a silently-dropped total-request-failure retry path). FCM is the one deliberate exception: it has no local emulator to test against, so dispatcher.fcm_test.go uses a hand-rolled fake FCMClient instead. Tests skip gracefully (t.Skip, not fail) when a backend isn't running locally, so go test . still works with no containers up — but the tests that matter for a given change need the real thing. See CLAUDE.md for the exact docker run commands (or make docker-up/make docker-down to start/stop them all at once) — these containers are shared across the whole gourdian25 workspace (grnoti, graudit, grcache, gourdiantoken all test against the same running instances, each using its own database/keyspace/DB-index).

Benchmarks: make bench (go test -bench=. -benchmem -benchtime=10s ./...) is a defined Makefile target, but there are currently no func Benchmark* functions anywhere in this repo — running it executes the regular test suite and prints no benchmark lines. Treat it as a target reserved for future use, not a populated benchmark suite today.

Error handling

Sentinel errors are used with errors.Is, never a per-error IsX(err) bool helper, matching every other gourdian repo. A backend-native error (mongo.ErrNoDocuments, pgx.ErrNoRows, redis.Nil, ...) is always translated to one of these before it can leak through a grnoti interface.

The 16 sentinels declared in errors.go:

Error Meaning
ErrClosed A method was called after Close
ErrBackendUnavailable A storage backend could not be reached
ErrInvalidEventID Event.EventID was empty
ErrNoTargetSpecified An Event has none of UserID, AnonymousID, or DeviceTokens set
ErrInvalidEventType Event.Type failed EventType.IsValid()
ErrInvalidPriority Event.Priority failed Priority.IsValid()
ErrTemplateNotFound No MessageTemplate registered for an event type, and no EventTypeCustom fallback either
ErrPreferencesNotFound No NotificationPreferences exist yet for a user (callers generally treat this as "use defaults")
ErrPreferencesUserIDRequired SavePreferences was called with an empty UserID
ErrExperimentNotFound No Experiment registered under the requested ID
ErrExperimentAlreadyExists CreateExperiment was called with an ID that already exists
ErrExperimentHasNoVariants An Experiment exists but has zero ExperimentVariants
ErrDLQEventNotFound No DLQEvent found for the requested event ID
ErrDLQEventNotClaimed MarkRetried was called for an event not currently in the claimed (retrying) state
ErrFCMClientNil A PushDispatcher was constructed with a nil FCM client
ErrFCMPayloadTooLarge A Message's estimated FCM payload size exceeds FCMMaxPayloadSize

Three more sentinels live alongside the feature they belong to rather than in errors.go, but follow the identical errors.Is convention: ErrCircuitOpen/ErrTooManyRequests (circuitbreaker.go) and ErrWorkerPoolFull (workerpool.go).

FCM failures additionally get a structured *FCMError (Code, Token, Message, wrapped Err), classified into a small FCMErrorCode enum (unspecified, invalid_argument, unregistered, sender_id_mismatch, quota_exceeded, unavailable, internal, third_party_auth_error) with IsRetryable()/IsPermanent() methods that drive real retry and invalid-token decisions in dispatcher.fcm.go — not just informational classification.

Limitations / out of scope

  • FCM has no local test emulator. Unlike Mongo/Postgres/Redis/Kafka, which are all tested against real local instances, FCM is tested with a hand-rolled fake FCMClient (dispatcher.fcm_test.go) — there is no equivalent "run it locally" option for Firebase Cloud Messaging.
  • CircuitBreaker and WorkerPool state is per-process, deliberately not centralized. Each service instance's breaker/queue is independent in-memory state — centralizing it across replicas was considered and rejected, since it risks a synchronized thundering-herd retry against FCM the moment it recovers, which is arguably worse than each instance recovering independently.
  • grevents publishing is always best-effort. A nil EventBus or a publish failure never affects the durable operation (dispatch/DLQ/idempotency) it follows — lifecycle events are an observability nicety, not part of the delivery guarantee.
  • DeviceToken is a bearer credential grnoti does not verify ownership of. grnoti does not check that a given token actually belongs to the UserID/AnonymousID it's registered under — that binding is entirely the caller's responsibility at TokenStore.SaveToken time.
  • No encryption at rest or in transit is provided by grnoti itself. TLS to Mongo/Postgres/Redis/Kafka, and encrypting any sensitive Event.Payload/Message.Data values before they reach grnoti, are the caller's responsibility.
  • Templates are not sanitized against injection. text/template (not html/template) renders Event.Payload values directly — a caller feeding untrusted input into a payload value is responsible for sanitizing it first.
  • Idempotency and DLQ keys are trusted, not authenticated. Any caller holding an IdempotencyStore/DLQHandler handle can mark arbitrary event IDs processed or replay dead-lettered events; there is no per-caller access control over these operations.
  • FCM credentials are never handled by grnoti. The PushDispatcher's FCM client is constructed and authenticated by the caller via the official Firebase Admin SDK — key management stays outside this library's scope.
  • Postgres schema management is additive only. Every New*Postgres* constructor applies CREATE TABLE/INDEX IF NOT EXISTS on connect; there is no down-migration, no versioning, and no support for evolving the schema beyond that. An ALTER TABLE, column type change, or backfill is entirely your own migration tool's job — set PostgresConfig.SkipSchemaEnsure: true once you own the schema that way. See docs/postgres.md.

See SECURITY.md for the complete scope-notes list this section draws from.

Development

make docker-up   # start the shared Postgres/Redis/Mongo/Kafka test containers
make precommit   # fmt + vet + lint + race + coverage-check
make docker-down # stop them when you're done

These containers are shared with graudit, grcache, and gourdiantoken (each gets its own database/keyspace) — see CLAUDE.md for backend setup, test scoping, and conventions.

Contributing

Issues and PRs are welcome at github.com/gourdian25/grnoti. Please run make precommit (fmt + vet + lint + race + coverage-check) before submitting.

License

MIT — see LICENSE.

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

View Source
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
)
View Source
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).

View Source
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"
)
View Source
const DefaultDLQCollection = "grnoti_dlq"

DefaultDLQCollection is the collection name used when MongoDLQHandlerConfig.CollectionName is empty.

View Source
const DefaultTokenCollection = "grnoti_tokens"

DefaultTokenCollection is the collection name used when MongoTokenStoreConfig.CollectionName is empty.

View Source
const FCMMaxPayloadSize = 4096

FCMMaxPayloadSize is FCM's documented maximum message payload size in bytes.

Variables

View Source
var (
	// ErrClosed indicates a method was called after Close.
	ErrClosed = errors.New("grnoti: closed")

	// ErrBackendUnavailable indicates a storage backend could not be
	// 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.

View Source
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.

View Source
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.

View Source
var ErrWorkerPoolFull = errors.New("grnoti: worker pool queue is full")

ErrWorkerPoolFull is returned by WorkerPool.Submit/SubmitAsync when the queue is full.

View Source
var Version = "v0.1.0"

Version is the semantic version of this module, matching its most recent git tag.

Functions

func FullJitterBackoff

func FullJitterBackoff(base, max time.Duration, attempt int) time.Duration

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 — but under a rare concurrent race on a brand-new (userID, experimentID) pair, more than one goroutine can independently observe "not yet assigned" and both proceed to assign+publish (the map/ cache write itself stays correct, since both computed the identical deterministic variant — see experiment.go's own doc comment). The result is at-least-once delivery for a given assignment, not exactly-once — an accepted characteristic of a best-effort side channel, matching 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.

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

func (e Event) GetTargetID() string

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

func (e Event) HasDirectTokens() bool

HasDirectTokens reports whether e targets an explicit set of device tokens, bypassing TokenStore lookup entirely.

func (Event) IsAnonymous

func (e Event) IsAnonymous() bool

IsAnonymous reports whether e targets an anonymous visitor (AnonymousID set, UserID not).

func (Event) IsAuthenticated

func (e Event) IsAuthenticated() bool

IsAuthenticated reports whether e targets a known user (UserID set).

func (Event) Validate

func (e Event) Validate() error

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

func (e EventType) IsValid() bool

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.

func (EventType) String

func (e EventType) String() string

String returns the underlying string value.

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.
	//
	// 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 NewCachedExperimentEngine) 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) Error

func (e *FCMError) Error() string

func (*FCMError) IsPermanent

func (e *FCMError) IsPermanent() bool

IsPermanent delegates to e.Code.IsPermanent.

func (*FCMError) IsRetryable

func (e *FCMError) IsRetryable() bool

IsRetryable delegates to e.Code.IsRetryable.

func (*FCMError) Unwrap

func (e *FCMError) Unwrap() error

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"
	FCMErrorCodeUnavailable       FCMErrorCode = "unavailable"
	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

func OrNop

func OrNop(l Logger) Logger

OrNop returns l if it is non-nil, otherwise NopLogger(). Every constructor in grnoti calls this once at construction time so every subsequent log call site can assume a non-nil Logger.

Parameters:

  • l: Logger — may be nil

Returns:

  • Logger: l unchanged if non-nil, otherwise NopLogger()

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"
)

func (Platform) IsValid

func (p Platform) IsValid() bool

IsValid reports whether p is one of the defined Platform constants.

func (Platform) String

func (p Platform) String() string

String returns the underlying string value.

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
	// SkipSchemaEnsure, if true, skips applying grnoti's embedded schema
	// on this connect call — for teams that manage the schema through
	// their own migration pipeline instead of grnoti's built-in
	// CREATE TABLE IF NOT EXISTS. See docs/postgres.md.
	SkipSchemaEnsure bool
	// 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"
)

func (Priority) IsValid

func (p Priority) IsValid() bool

IsValid reports whether p is one of the defined Priority constants.

func (Priority) String

func (p Priority) String() string

String returns the underlying string value.

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.

type WorkerPoolStats

type WorkerPoolStats struct {
	Workers      int
	QueueSize    int
	QueuedEvents int
	QueueUsage   float64 // QueuedEvents / QueueSize
}

WorkerPoolStats is a point-in-time snapshot of a WorkerPool's queue.

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

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL