delivery

package
v1.1.12-test Latest Latest
Warning

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

Go to latest
Published: Aug 28, 2026 License: MIT Imports: 18 Imported by: 0

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

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

type Delivery interface {
	GetID() string
	DeliveryState() *State
}

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

type Engine[D Delivery] struct {
	*negotiate.SyncJob[D]
	// contains filtered or unexported fields
}

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

func NewEngine[D Delivery](store Store[D], sender Sender[D], cfg Config) *Engine[D]

NewEngine wires the discovery loop and (optional) consumer around a Store and Sender.

func (*Engine[D]) Enqueue

func (e *Engine[D]) Enqueue(ids ...string) error

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

func (e *Engine[D]) Process(ctx context.Context, id string) error

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.

func (*Engine[D]) Start

func (e *Engine[D]) Start() <-chan error

Start runs the Pub/Sub consumer (when a queue is configured) and the leader-elected discovery loop.

func (*Engine[D]) Stop

func (e *Engine[D]) Stop()

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 NewGormStore[D Delivery](dbRw, dbRo *gorm.DB, newRecord func() D) *GormStore[D]

func (*GormStore[D]) GetByID

func (g *GormStore[D]) GetByID(ctx context.Context, id string) (D, error)

func (*GormStore[D]) ListDue

func (g *GormStore[D]) ListDue(ctx context.Context, now time.Time) ([]D, error)

func (*GormStore[D]) MarkPushed

func (g *GormStore[D]) MarkPushed(ctx context.Context, id string) error

func (*GormStore[D]) Save

func (g *GormStore[D]) Save(ctx context.Context, d D) error

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

type Sender[D Delivery] interface {
	Send(ctx context.Context, d D) error
}

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)

Jump to

Keyboard shortcuts

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