Documentation
¶
Overview ¶
Package delivery is a reusable engine for durable, idempotent, retried, audited, Pub/Sub-backed delivery pipelines. A pipeline supplies a record model (embedding State), a Store (persistence) and a Sender (the actual delivery); the Engine owns the leader-elected discovery loop, the Pub/Sub consumer, and the claim -> send -> finalize lifecycle with retry/backoff and an error-history audit trail. Record creation/enqueue stays pipeline-specific (domain logic).
Index ¶
Constants ¶
const ( StatusPending = "pending" StatusInProgress = "in_progress" StatusDelivered = "delivered" StatusFailed = "failed" StatusPermanentlyFailed = "permanently_failed" )
Delivery statuses. A record advances pending -> in_progress -> delivered, or in_progress -> failed (a retry is scheduled) -> ... -> permanently_failed.
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type ClassifiedError ¶
type ClassifiedError struct {
OriginalError error
ErrorType ErrorType
RetryAfter *time.Duration // for rate-limit errors
IsRetryable bool
Message string
}
ClassifiedError wraps an error with its classification.
func ClassifyError ¶
func ClassifyError(err error, deliveryContext string) *ClassifiedError
ClassifyError classifies an error to determine retry behavior. deliveryContext is an optional label (e.g. "email", "slack") used only for diagnostics.
Structured signals win: an HTTP status carried by the error, a DomainError kind, or a known sentinel is authoritative. Message text is only consulted for errors that carry no structure, since matching prose is inherently ambiguous. Anything unrecognized stays retryable — dropping a delivery is worse than sending it twice.
func (*ClassifiedError) Error ¶
func (e *ClassifiedError) Error() string
type Config ¶
type Config struct {
// Name uniquely identifies the pipeline (leader-election class + job/service
// name + log label).
Name string
// SyncInterval is how often the discovery loop scans for due records.
SyncInterval time.Duration
// MaxRetries caps retry attempts (0 = strategy default).
MaxRetries int
// Elector gates the discovery loop to a single leader.
Elector leaderElection.Elector
// Queue, when non-nil, distributes processing via Pub/Sub. Nil = poll-only
// (the leader processes inline).
Queue QueueClient
// OnPermanentFailure, when set, is called after a record exhausts retries
// (e.g. to raise an alert). Optional.
OnPermanentFailure func(d Delivery, err error)
}
Config configures an Engine.
type Delivery ¶
Delivery is implemented by each pipeline's record model so the engine can drive the shared lifecycle while the model carries its own payload columns.
type Engine ¶
Engine drives the shared delivery lifecycle. It satisfies the Name/Start/Stop background-job shape so callers can register it like any other service.
func NewEngine ¶
NewEngine wires the discovery loop and (optional) consumer around a Store and Sender.
func (*Engine[D]) Enqueue ¶
Enqueue is a best-effort, low-latency nudge that submits already-committed outbox rows into the delivery system immediately, rather than waiting for the leader's discovery loop to find them. It is NOT the durability guarantee — the leader's discovery loop re-pushes any row still lacking pushed_at — so callers may ignore its error. No-op when the queue is disabled, since processing on a non-leader pod could duplicate sends (the leader poll picks the rows up).
func (*Engine[D]) Process ¶
Process drives a single delivery through claim -> send -> finalize. Idempotent: a record already delivered or permanently failed is skipped, and the claim (status=in_progress) guards against concurrent processing.
type ErrorEntry ¶
type ErrorEntry struct {
AttemptTimestamp string `json:"attemptTimestamp"`
Message string `json:"message"`
// Detail is the underlying error exactly as the Sender returned it. Message is
// only a classification summary and is deliberately generic ("Authentication or
// authorization failed"); Detail is what actually went wrong, and is the field
// to read when diagnosing a failed delivery.
Detail string `json:"detail,omitempty"`
AttemptNumber int `json:"attemptNumber"`
ErrorType string `json:"errorType"`
}
ErrorEntry is one entry in a delivery's structured error history.
type ErrorType ¶
type ErrorType string
ErrorType classifies a delivery error to determine retry behavior.
const ( ErrorTypeTransient ErrorType = "transient" // network issues, timeouts, 5xx ErrorTypeRateLimit ErrorType = "rate_limit" // 429 ErrorTypeAuth ErrorType = "auth" // 401/403 — credentials invalid ErrorTypePermanent ErrorType = "permanent" // 400/404 — won't succeed on retry ErrorTypeUnknown ErrorType = "unknown" // unclassified (retried conservatively) )
type GormStore ¶
type GormStore[D Delivery] struct { // contains filtered or unexported fields }
GormStore is a generic gorm-backed Store over any model embedding State. A new pipeline gets persistence for free by passing a newRecord factory:
NewGormStore(dbRw, dbRo, func() *Foo { return &Foo{} })
func NewGormStore ¶
func (*GormStore[D]) MarkPushed ¶
type QueueClient ¶
type QueueClient = pubsub.Client[TaskMessage]
QueueClient is the per-pipeline Pub/Sub client (its own topic/subscription/DLQ).
func NewQueueClient ¶
func NewQueueClient(ctx context.Context, cfg pubsub.ClientConfig) (QueueClient, error)
NewQueueClient creates a delivery queue client (auto-creates topic/DLQ/sub).
type RetryStrategy ¶
type RetryStrategy struct {
MaxRetries int
BaseDelay time.Duration
MaxDelay time.Duration
Multiplier float64
JitterFraction float64 // 0.0..1.0
}
RetryStrategy determines when to retry based on error type and attempt count.
func DefaultRetryStrategy ¶
func DefaultRetryStrategy() *RetryStrategy
DefaultRetryStrategy returns a sensible default (exponential backoff + jitter).
func (*RetryStrategy) ApplyRetry ¶
func (s *RetryStrategy) ApplyRetry(state *State, err error) bool
ApplyRetry classifies a send error and updates the delivery State for the next attempt: a retryable error sets status=failed with a scheduled NextRetryAt; otherwise status=permanently_failed. Either way the error is appended to the audit history. Returns true if a retry was scheduled.
func (*RetryStrategy) CalculateNextRetry ¶
func (s *RetryStrategy) CalculateNextRetry(classifiedErr *ClassifiedError, attemptCount int) (time.Time, bool)
CalculateNextRetry returns when to retry next and whether to retry at all.
type Sender ¶
Sender performs the actual delivery for a record. Errors are returned verbatim so the engine can classify them for retry (include HTTP status text where possible, e.g. "sendgrid returned 429: ...").
type State ¶
type State struct {
Status string `gorm:"type:text;not null;default:'pending';column:status" json:"status"`
AttemptCount int `gorm:"type:integer;not null;default:0;column:attempt_count" json:"attemptCount"`
// PushedAt records the first time this record was submitted into the delivery
// system (published to the queue, or handed to the leader's inline processor
// in poll-only mode). It is nil for a freshly-created outbox row that the
// initiating action committed but that has not yet reached the delivery
// system. The leader-elected discovery loop is responsible for pushing every
// unpushed row and stamping this column, so a crash between the initiating
// commit and the push can never lose the message.
PushedAt *time.Time `gorm:"type:timestamptz;column:pushed_at" json:"pushedAt,omitempty"`
DeliveredAt *time.Time `gorm:"type:timestamptz;column:delivered_at" json:"deliveredAt,omitempty"`
NextRetryAt *time.Time `gorm:"type:timestamptz;column:next_retry_at" json:"nextRetryAt,omitempty"`
LastAttemptAt *time.Time `gorm:"type:timestamptz;column:last_attempt_at" json:"lastAttemptAt,omitempty"`
LastErrorType *string `gorm:"type:text;column:last_error_type" json:"lastErrorType,omitempty"`
ErrorMessages []ErrorEntry `gorm:"type:jsonb;column:error_messages;serializer:json" json:"errorMessages,omitempty"`
CreatedAt time.Time `gorm:"type:timestamptz;not null;default:now();column:created_at" json:"createdAt"`
UpdatedAt time.Time `gorm:"type:timestamptz;not null;default:now();column:updated_at" json:"updatedAt"`
}
State holds the delivery-lifecycle fields the engine manages. Embed it ANONYMOUSLY in a pipeline's record model so gorm flattens the columns and JSON promotes the fields to the parent object (preserving existing API/UI shapes):
type Foo struct {
Id string `gorm:"primaryKey;column:id"`
delivery.State
// ... pipeline-specific payload columns ...
}
func (f *Foo) GetID() string { return f.Id }
func (f *Foo) DeliveryState() *delivery.State { return &f.State }
type Store ¶
type Store[D Delivery] interface { GetByID(ctx context.Context, id string) (D, error) Save(ctx context.Context, d D) error // MarkPushed stamps pushed_at (once) to record that the record has been // submitted into the delivery system. It is a targeted column update, NOT a // full Save, so it never clobbers a concurrent consumer's status write; and // it only sets the column when still null, so retry re-pushes are no-ops. MarkPushed(ctx context.Context, id string) error // ListDue returns records that need pushing or (re)processing as of now: // unpushed outbox rows, due retries, pushed-but-never-consumed rows, and // stuck in_progress rows. ListDue(ctx context.Context, now time.Time) ([]D, error) }
Store persists delivery records. The engine only needs these operations; pipeline-specific creation/listing lives in the pipeline's own repository.
type TaskMessage ¶
type TaskMessage struct {
DeliveryID string `json:"deliveryId"`
PublishedAt time.Time `json:"publishedAt"`
}
TaskMessage is the Pub/Sub payload — just the delivery id; the engine re-loads the record. One message shape across all pipelines (each uses its own topic).
func (TaskMessage) Marshal ¶
func (m TaskMessage) Marshal() ([]byte, error)