eventwake

package
v0.1.59 Latest Latest
Warning

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

Go to latest
Published: Sep 14, 2026 License: Apache-2.0 Imports: 13 Imported by: 0

Documentation

Index

Constants

View Source
const (
	EnvelopeVersion = 1
)

Variables

View Source
var ErrSubscriptionClosed = errors.New("event wake subscription is closed")
View Source
var ErrUnknownTopic = errors.New("event wake topic is not configured")
View Source
var ErrWakeSourceDegraded = errors.New("event wake source is degraded")

Functions

func RunScheduler

func RunScheduler(
	ctx context.Context,
	subscription *Subscription,
	config SchedulerConfig,
	run SchedulerRun,
) error

RunScheduler replaces fixed-frequency claim attempts with three bounded wake sources: a transactional topic notification, the worker's earliest durable due time, and a low-frequency reconciliation timer. The callback remains solely responsible for the authoritative PostgreSQL claim/CAS.

Types

type Envelope

type Envelope struct {
	Version    int       `json:"version"`
	Topic      string    `json:"topic"`
	ResourceID string    `json:"resource_id"`
	Generation uint64    `json:"generation"`
	ProducedAt time.Time `json:"produced_at"`
}

func ParseEnvelope

func ParseEnvelope(encoded []byte) (Envelope, error)

type Hub

type Hub struct {
	// contains filtered or unexported fields
}

Hub is a process-local, non-authoritative wake broadcaster. It retains no business payload and intentionally drops wakes for keys with no subscribers. Consumers must subscribe before querying PostgreSQL.

func NewHub

func NewHub() *Hub

func (*Hub) Publish

func (h *Hub) Publish(key string)

Publish advances the key generation and broadcasts without waiting for any consumer. Repeated publishes safely coalesce at the generation boundary.

func (*Hub) PublishAll

func (h *Hub) PublishAll()

PublishAll is used after listener recovery so every local waiter re-reads PostgreSQL and closes any gap left by non-durable notifications.

func (*Hub) Subscribe

func (h *Hub) Subscribe(key string) *Subscription

type Infrastructure

type Infrastructure struct {
	// contains filtered or unexported fields
}

Infrastructure owns one PostgreSQL LISTEN connection and a bounded local Router. It is advisory only and does not participate in readiness or replace existing worker scheduling until an explicit later cutover.

func NewPostgresInfrastructure

func NewPostgresInfrastructure(
	pool *pgxpool.Pool,
	channels []string,
	topics []string,
) (*Infrastructure, error)

func (*Infrastructure) Health

func (i *Infrastructure) Health() ListenerHealth

func (*Infrastructure) ListenerStats

func (i *Infrastructure) ListenerStats() ListenerStats

func (*Infrastructure) Run

func (i *Infrastructure) Run(ctx context.Context) error

func (*Infrastructure) Subscribe

func (i *Infrastructure) Subscribe(topic, resourceID string) (*Subscription, error)

func (*Infrastructure) SubscribeTopic

func (i *Infrastructure) SubscribeTopic(topic string) (*Subscription, error)

func (*Infrastructure) TopicStats

func (i *Infrastructure) TopicStats() map[string]TopicStats

type Listener

type Listener struct {
	// contains filtered or unexported fields
}

func NewPostgresListener

func NewPostgresListener(pool *pgxpool.Pool, config ListenerConfig) (*Listener, error)

func (*Listener) Health

func (l *Listener) Health() ListenerHealth

func (*Listener) Run

func (l *Listener) Run(ctx context.Context) error

func (*Listener) Stats

func (l *Listener) Stats() ListenerStats

type ListenerConfig

type ListenerConfig struct {
	Channels   []string
	Topics     []string
	MinBackoff time.Duration
	MaxBackoff time.Duration
	Dispatch   func(context.Context, Envelope)
	OnRecovery func(uint64)
}

type ListenerHealth

type ListenerHealth struct {
	Connected  bool
	Generation uint64
	Reason     string
}

type ListenerStats

type ListenerStats struct {
	Accepted        uint64
	Rejected        uint64
	RejectedReasons map[string]uint64
	ConnectFailures uint64
	ListenFailures  uint64
	WaitFailures    uint64
	Reconnects      uint64
}

type Notification

type Notification struct {
	Channel string
	Payload string
}

type Router

type Router struct {
	// contains filtered or unexported fields
}

Router projects advisory envelopes into independent process-local Hubs. It records only bounded counters and timestamps for shadow comparison; it never stores business payloads or reads PostgreSQL.

func NewRouter

func NewRouter(topics []string) (*Router, error)

func (*Router) Dispatch

func (r *Router) Dispatch(_ context.Context, envelope Envelope)

Dispatch is safe to call from Listener: Hub.Publish is non-blocking and the stats update is bounded. The awakened consumer remains responsible for querying the authoritative PostgreSQL state.

func (*Router) Recover

func (r *Router) Recover(_ uint64)

Recover broadcasts to every local waiter after a LISTEN (re)connect. Since NOTIFY is not durable, each awakened consumer must re-read PostgreSQL.

func (*Router) Stats

func (r *Router) Stats() map[string]TopicStats

func (*Router) Subscribe

func (r *Router) Subscribe(topic, resourceID string) (*Subscription, error)

func (*Router) SubscribeTopic

func (r *Router) SubscribeTopic(topic string) (*Subscription, error)

SubscribeTopic coalesces every resource notification for one work type. It is intended for process-level workers; resource waiters must continue to use Subscribe so unrelated Runs do not wake one another.

type SchedulerConfig

type SchedulerConfig struct {
	ReconcileInterval time.Duration
	ErrorRetry        time.Duration
	MinimumDelay      time.Duration
	HealthCheck       time.Duration
	Healthy           func() bool
}

type SchedulerResult

type SchedulerResult struct {
	NextDelay time.Duration
	HasNext   bool
}

type SchedulerRun

type SchedulerRun func(context.Context, string) (SchedulerResult, error)

type Subscription

type Subscription struct {
	// contains filtered or unexported fields
}

func (*Subscription) Close

func (s *Subscription) Close()

func (*Subscription) Generation

func (s *Subscription) Generation() uint64

func (*Subscription) Wait

func (s *Subscription) Wait(ctx context.Context) (uint64, error)

type TopicSource

type TopicSource interface {
	Health() ListenerHealth
	SubscribeTopic(string) (*Subscription, error)
}

TopicSource is the deliberately small worker-facing surface. Notifications are advisory: Health controls cutover/fallback, and PostgreSQL claims remain the only authoritative ownership transition.

type TopicStats

type TopicStats struct {
	Accepted       uint64
	RecoveryWakes  uint64
	LastGeneration uint64
	LastWakeAt     time.Time
	LastWakeLag    time.Duration
	MaxWakeLag     time.Duration
}

Jump to

Keyboard shortcuts

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