Documentation
¶
Index ¶
- Constants
- Variables
- func RunScheduler(ctx context.Context, subscription *Subscription, config SchedulerConfig, ...) error
- type Envelope
- type Hub
- type Infrastructure
- func (i *Infrastructure) Health() ListenerHealth
- func (i *Infrastructure) ListenerStats() ListenerStats
- func (i *Infrastructure) Run(ctx context.Context) error
- func (i *Infrastructure) Subscribe(topic, resourceID string) (*Subscription, error)
- func (i *Infrastructure) SubscribeTopic(topic string) (*Subscription, error)
- func (i *Infrastructure) TopicStats() map[string]TopicStats
- type Listener
- type ListenerConfig
- type ListenerHealth
- type ListenerStats
- type Notification
- type Router
- type SchedulerConfig
- type SchedulerResult
- type SchedulerRun
- type Subscription
- type TopicSource
- type TopicStats
Constants ¶
const (
EnvelopeVersion = 1
)
Variables ¶
var ErrSubscriptionClosed = errors.New("event wake subscription is closed")
var ErrUnknownTopic = errors.New("event wake topic is not configured")
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 ¶
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 (*Hub) Publish ¶
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 (*Infrastructure) Health ¶
func (i *Infrastructure) Health() ListenerHealth
func (*Infrastructure) ListenerStats ¶
func (i *Infrastructure) ListenerStats() ListenerStats
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) Stats ¶
func (l *Listener) Stats() ListenerStats
type ListenerConfig ¶
type ListenerHealth ¶
type ListenerStats ¶
type Notification ¶
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 (*Router) Dispatch ¶
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 ¶
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 SchedulerResult ¶
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
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.