Documentation
¶
Overview ¶
Package bgloop is the one way the platform runs work in the background (#1897). Every periodic sweep, poll/wake worker, heartbeat and event consumer goes through Run or Consume, which is what makes a loop observable without its author remembering to: each iteration is counted by result (background_loop_iterations_total{loop,result}), timed (background_loop_duration_seconds{loop}) and, when it succeeds, stamped (background_loop_last_success_timestamp_seconds{loop}), so a loop that has stopped reads as a timestamp that stopped advancing rather than as silence.
Each iteration also gets a span: a root span by default, since a background iteration has no caller to continue, or a child of the caller's span for a loop that runs inside a unit of work (a lease heartbeat). The iteration's context carries that span, so a log record the body writes with it is correlated with the iteration's trace. A loop whose iterations are mostly idle polls (a queue worker) sets SpanPerUnit and opens a span per unit of work it finds with Unit instead, keeping a trace per claimed job rather than one per empty poll.
.semgrep/go-background-loop.yml refuses time.NewTicker, time.Tick and a for/select loop anywhere else, so a new loop cannot ship without these signals.
Index ¶
- Constants
- Variables
- func Consume[T any](ctx context.Context, e Events[T])
- func Metrics() *observability.Metrics
- func Purged(ctx context.Context, name string, n int64)
- func Result(err error) string
- func Run(ctx context.Context, l Loop)
- func SetDefaultMetrics(m *observability.Metrics)
- func Stopped(ctx context.Context, stop <-chan struct{}) bool
- func Unit(ctx context.Context, name string, fn func(ctx context.Context) error) error
- func Woken(ctx context.Context) bool
- type Events
- type Listen
- type Loop
Constants ¶
const ( // Request-path state the platform keeps in memory, swept on a timer. NameSessionGateCleanup = "session_gate_cleanup" NameSessionEnrichmentCleanup = "session_enrichment_cleanup" NameSessionErrorsCleanup = "session_errors_cleanup" NameSearchGateCleanup = "search_gate_cleanup" NamePKCECleanup = "pkce_cleanup" NameSessionCleanup = "session_cleanup" NameOAuthCleanup = "oauth_cleanup" NameOAuthStoreCleanup = "oauth_store_cleanup" // Retention and maintenance sweeps. NameAuditMaintenance = "audit_maintenance" NameAuthEventsPrune = "authevents_prune" NameCallCatalogSweep = "call_catalog_sweep" NameIndexJobRetention = "indexjob_retention" NameRetention = "retention" // Queue workers and their sub-loops. NameIndexJobWorker = "indexjob_worker" NameIndexJobHeartbeat = "indexjob_heartbeat" NameIndexJobReaper = "indexjob_reaper" NameIndexJobReconciler = "indexjob_reconciler" NameIndexJob = "indexjob" NameNotifyWorker = "notification_worker" NameNotifyDelivery = "notification_delivery" NameScriptScheduler = "script_scheduler" NameScriptWorker = "script_worker" NameScriptShed = "script_shed" NameScriptRun = "script_run" NameScriptProgress = "script_progress" NameThumbnailWorker = "thumbnail_worker" NameThumbnailRender = "thumbnail_render" NameWebhookCompactor = "webhook_compactor" NameWebhookSegments = "webhook_segments" NameWebhookSegmentWrite = "webhook_segment_write" NameWebhookCompaction = "webhook_compaction" NameWebhookRetention = "webhook_retention" NameWebhookRefresh = "webhook_source_refresh" NameWebhookStats = "webhook_stats_flush" // Alerts, keepalives and watchers. NameConnOAuthRefresh = "connection_oauth_refresh" NameRevocationEscalate = "connection_revocation_escalator" NameReviewQueueAlert = "review_queue_alert" NameMemoryStaleness = "memory_staleness" NameReloadBus = "reload_bus" // LISTEN connections (Listen). The two pglisten channels report under // "listen_" + their channel name. NameListenIndexJobs = "listen_index_jobs" NameListenSessions = "listen_session_broadcast" ListenPrefix = "listen_" )
Loop names. Every loop the platform runs is named here, which is what bounds the loop label on the background series: a name is added beside the loop that uses it, never built from data. Two additions are derived from names here and stay bounded with them: a retention sweep reports under Retention + "_" + its sweep name (internal/platform/retention declares the sweeps), and a LISTEN connection under Listen* below.
Variables ¶
var ErrSkipped = errors.New("bgloop: work held elsewhere")
ErrSkipped is returned by an iteration that found its work held elsewhere (another replica's advisory lock) and did nothing. It is counted as result "skipped", is not an error, and does not advance the last-success time.
var ErrStop = errors.New("bgloop: stop")
ErrStop is returned by an iteration to end its loop: a heartbeat whose lease went to another worker. The iteration is counted as ok.
Functions ¶
func Metrics ¶
func Metrics() *observability.Metrics
Metrics returns the installed recorder (nil when none is), for a loop body recording its own series beside the loop's.
func Run ¶
Run runs l until ctx ends or l.Stop is closed. It blocks; callers start it on their own goroutine.
func SetDefaultMetrics ¶
func SetDefaultMetrics(m *observability.Metrics)
SetDefaultMetrics installs the recorder loops record through. Nil turns recording off.
func Stopped ¶
Stopped reports whether stop is closed or ctx has ended, without blocking: the check a drain loop makes between units of work.
Types ¶
type Events ¶
type Events[T any] struct { // Name labels the loop's series and spans (names.go). Name string // C is the channel consumed. The loop ends when it is closed. C <-chan T // Stop ends the loop, as ctx ending does. Optional. Stop <-chan struct{} // SpanPerUnit opens no span per event; for a high-rate, trivial event // such as a LISTEN wakeup, whose work is traced where it is done. SpanPerUnit bool // Body handles one event. Body func(ctx context.Context, ev T) error }
Events describes a loop that consumes a channel of events.
type Listen ¶
type Listen struct {
// contains filtered or unexported fields
}
Listen reports one LISTEN connection's state on the pg_listen_* series: whether it is connected, how often it reconnected, and how long since it last delivered a notification (#1897). Each adapter that opens a pq.Listener feeds its event callback and its notifications through one.
A half-open connection is not detected by a ping: pq.Listener.Ping holds the listener's lock for the round trip, so a ping on a connection whose peer is gone would hold Close until the operating system gave up on the socket. The last-notification age is the signal for a connection that stays "connected" and delivers nothing. A nil *Listen records nothing.
func NewListen ¶
NewListen returns the state reporter for the LISTEN connection named name (a Listen* constant in names.go).
func (*Listen) Event ¶
func (l *Listen) Event(ev pq.ListenerEventType)
Event records a pq.Listener lifecycle event: connected and reconnected set the gauge to 1 (a reconnect also counts one), a drop or a failed attempt sets it to 0.
type Loop ¶
type Loop struct {
// Name labels every series and span the loop produces. It is one of the
// Name* constants in names.go.
Name string
// Every is the wait between iterations. Ignored when Next is set.
Every time.Duration
// Next, when set, returns the wait before the next iteration, read after
// each one: a loop that found work runs again at once, one that found a
// lock held retries sooner.
Next func() time.Duration
// Immediate runs the first iteration at start instead of one wait in.
Immediate bool
// Wake starts an iteration early: a LISTEN notification, an enqueue on
// this replica, an administrator's request.
Wake <-chan struct{}
// Stop ends the loop, as ctx ending does. Optional.
Stop <-chan struct{}
// Child makes each iteration's span a child of ctx's span instead of a
// root: a loop running inside a unit of work, such as a lease heartbeat.
Child bool
// SpanPerUnit opens no span for the iteration; the body opens one per
// unit of work it finds with Unit.
SpanPerUnit bool
// Body is one iteration. It returns ErrSkipped when the work was held
// elsewhere and an error when it failed; the loop continues either way.
Body func(ctx context.Context) error
// Final runs once after the loop stops, with a context that is not
// canceled: a buffer writing what it still holds.
Final func(ctx context.Context)
}
Loop describes one background loop.