From 3f8a7b670778a513c7d6d68de50de6f277e18a0f Mon Sep 17 00:00:00 2001 From: Aloento <11802769+Aloento@users.noreply.github.com> Date: Sun, 4 Oct 2026 18:40:25 +0200 Subject: [PATCH 1/2] Make the checker safe for multi-replica deployments - Drop the in-memory scan cursors (lastMntID/lastInfoID) and the active-event bookkeeping; every round now scans from the start, with idempotency coming from the existing status derivation. - Remove the now-unused after parameter from GetMaintenances, GetInfoEvents and getEventsByType. - Add db.WithAdvisoryLock: a non-blocking session-level advisory lock acquired and released on the same dedicated connection, returning the new ErrLockBusy sentinel when the lock is held. - Wrap the whole scan round in the advisory lock (key 9001, SD3 reserved range); a busy lock skips the round with a debug log. - The checker now reuses the app's database pool and notification publisher instead of opening its own pool and publisher; the manual wiring in main is removed. - Fix the checker shutdown race: Run/Shutdown now use a cancel function and a done channel, so repeated shutdowns no longer panic. --- cmd/main.go | 25 +------ internal/app/app.go | 8 ++- internal/checker/checker.go | 78 +++++++++++++-------- internal/checker/info.go | 27 +------- internal/checker/maintenance.go | 36 +--------- internal/db/errors.go | 1 + internal/db/event_types.go | 6 +- internal/db/info.go | 4 +- internal/db/lock.go | 38 +++++++++++ internal/db/maintenances.go | 4 +- tests/advisory_lock_test.go | 102 ++++++++++++++++++++++++++++ tests/checker_notifications_test.go | 24 +++++-- 12 files changed, 227 insertions(+), 126 deletions(-) create mode 100644 internal/db/lock.go create mode 100644 tests/advisory_lock_test.go diff --git a/cmd/main.go b/cmd/main.go index a293bea..bad46f2 100644 --- a/cmd/main.go +++ b/cmd/main.go @@ -29,22 +29,7 @@ func main() { logger.Fatal("fail to init app", zap.Error(err)) } - ch, err := checker.New(c, logger) - if err != nil { - logger.Error("fail to init checker", zap.Error(err)) - } - - // Wire the checker's publisher to the app's single delivery worker so - // checker-driven transitions wake it immediately (same shared queue). - if ch != nil { - // NotifyFunc is nil when notifications are disabled; leave the publisher - // unwired rather than installing a nil callback. - if notify := s.NotifyFunc(); notify != nil { - ch.Publisher().SetNotify(notify) - } - } - - stopCh := make(chan struct{}) + ch := checker.New(s.DB, logger, s.Publisher()) ctx, done := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM) defer done() @@ -55,9 +40,7 @@ func main() { } }() - go func() { - ch.Run(stopCh) - }() + go ch.Run() <-ctx.Done() s.Log.Info("shutdown app") @@ -66,9 +49,7 @@ func main() { logger.Fatal("app shutdown failed", zap.Error(err)) } - if err = ch.Shutdown(stopCh); err != nil { - logger.Fatal("checker shutdown failed", zap.Error(err)) - } + ch.Shutdown() logger.Info("app exited") } diff --git a/internal/app/app.go b/internal/app/app.go index d7f4c2a..a26ced9 100644 --- a/internal/app/app.go +++ b/internal/app/app.go @@ -120,7 +120,7 @@ func buildWorker( } // NotifyFunc returns the worker's wake-up callback, or nil when notifications are -// disabled. Used to wire the checker's publisher to the same worker. +// disabled. func (a *App) NotifyFunc() func() { if a.worker == nil { return nil @@ -128,6 +128,12 @@ func (a *App) NotifyFunc() func() { return a.worker.Notify } +// Publisher returns the app's notification publisher, already wired to the +// delivery worker's Notify. +func (a *App) Publisher() *notification.Publisher { + return a.api.Publisher() +} + func (a *App) Run() error { if a.worker != nil { var ctx context.Context diff --git a/internal/checker/checker.go b/internal/checker/checker.go index 6579f7a..c25583d 100644 --- a/internal/checker/checker.go +++ b/internal/checker/checker.go @@ -1,40 +1,57 @@ package checker import ( + "context" + "errors" "sync" "time" "go.uber.org/zap" - "github.com/stackmon/otc-status-dashboard/internal/conf" "github.com/stackmon/otc-status-dashboard/internal/db" "github.com/stackmon/otc-status-dashboard/internal/notification" ) const defaultPeriod = time.Minute * 2 +// scanLockKey guards the full scan across replicas. It lives in the SD3 +// reserved advisory-lock range 9000-9099; it will move to internal/scheduler +// when the unified scheduler lands. +const scanLockKey int64 = 9001 + type Checker struct { db *db.DB log *zap.Logger notifier *notification.Publisher - // lastIDs are the earliest planned or in_progress maintenance/info events ID. - lastMntID uint - lastInfoID uint + mu sync.Mutex + cancel context.CancelFunc + done chan struct{} } -func New(c *conf.Config, log *zap.Logger) (*Checker, error) { - dbNew, err := db.New(c) - if err != nil { - return nil, err +// New builds a Checker on the app's shared database pool and notification +// publisher. It owns neither: the pool and the publisher are closed and +// wired by the app. +func New(database *db.DB, log *zap.Logger, notifier *notification.Publisher) *Checker { + return &Checker{db: database, log: log, notifier: notifier, done: make(chan struct{})} +} + +func (ch *Checker) Check() { + // One lock per round so only one replica scans at a time; the scan is + // idempotent, so a skipped round costs nothing. + err := ch.db.WithAdvisoryLock(context.Background(), scanLockKey, func(context.Context) error { + ch.runScan() + return nil + }) + if errors.Is(err, db.ErrLockBusy) { + ch.log.Debug("another replica holds the scan lock, skipping this round") + return } - ncfg, err := notification.ConfigFromConf(c) if err != nil { - return nil, err + ch.log.Error("failed to acquire the scan lock", zap.Error(err)) } - return &Checker{db: dbNew, log: log, notifier: notification.NewPublisher(ncfg, dbNew)}, nil } -func (ch *Checker) Check() { +func (ch *Checker) runScan() { var wg sync.WaitGroup wg.Add(1) @@ -58,14 +75,21 @@ func (ch *Checker) Check() { wg.Wait() } -func (ch *Checker) Run(done chan struct{}) { +func (ch *Checker) Run() { ch.log.Info("checker is started") + ctx, cancel := context.WithCancel(context.Background()) + ch.mu.Lock() + ch.cancel = cancel + ch.mu.Unlock() + defer cancel() + ticker := time.NewTicker(defaultPeriod) defer ticker.Stop() for { //nolint:nolintlint select { - case <-done: + case <-ctx.Done(): + close(ch.done) return case <-ticker.C: ch.Check() @@ -73,20 +97,16 @@ func (ch *Checker) Run(done chan struct{}) { } } -func (ch *Checker) Shutdown(done chan struct{}) error { +// Shutdown stops the Run loop and waits for it to exit. It is safe to call +// multiple times and without Run having started. +func (ch *Checker) Shutdown() { ch.log.Info("start to shutdown checker") - done <- struct{}{} - close(done) - return ch.db.Close() -} - -// Close releases the checker's database pool without going through the Run loop. -func (ch *Checker) Close() error { - return ch.db.Close() -} - -// Publisher returns the checker's notification publisher so the delivery worker's -// Notify can be wired in during app startup. -func (ch *Checker) Publisher() *notification.Publisher { - return ch.notifier + ch.mu.Lock() + cancel := ch.cancel + ch.mu.Unlock() + if cancel == nil { + return + } + cancel() + <-ch.done } diff --git a/internal/checker/info.go b/internal/checker/info.go index 345cb53..a58c613 100644 --- a/internal/checker/info.go +++ b/internal/checker/info.go @@ -1,7 +1,6 @@ package checker import ( - "slices" "time" "go.uber.org/zap" @@ -48,16 +47,12 @@ func (st *InfoStatusHistory) setStatus(status event.Status) { func (ch *Checker) CheckInfoEvents() error { ch.log.Info("check info event statuses") - if ch.lastInfoID == 0 { - ch.log.Info("no last completed info event, starting from the beginning") - } - infos, err := ch.db.GetInfoEvents(ch.lastInfoID) + infos, err := ch.db.GetInfoEvents() if err != nil { return err } - var activeInfoEvents []uint for _, info := range infos { sHistory := calculateInfoStatusHistory(info) actualStatus := calculateCurrentInfoStatus(sHistory, info) @@ -65,10 +60,8 @@ func (ch *Checker) CheckInfoEvents() error { switch actualStatus { case event.InfoPlanned: ch.fixInfoMissedStatuses(event.InfoPlanned, sHistory, info) - activeInfoEvents = append(activeInfoEvents, info.ID) case event.InfoActive: ch.fixInfoMissedStatuses(event.InfoActive, sHistory, info) - activeInfoEvents = append(activeInfoEvents, info.ID) case event.InfoCompleted: ch.fixInfoMissedStatuses(event.InfoCompleted, sHistory, info) case event.InfoCancelled: @@ -86,24 +79,6 @@ func (ch *Checker) CheckInfoEvents() error { } } - if len(activeInfoEvents) == 0 { - for _, mn := range infos { - if mn.ID > ch.lastInfoID { - ch.lastInfoID = mn.ID - } - } - ch.log.Debug( - "there are no actual info events, set the last ID to the last one", - zap.Uint("lastInfoID", ch.lastInfoID), - ) - } else { - ch.lastInfoID = slices.Min(activeInfoEvents) - ch.log.Debug( - "set the last ID to the earliest planned or in_progress info event", - zap.Uint("lastInfoID", ch.lastInfoID), - ) - } - ch.log.Info("finished checking info events") return nil diff --git a/internal/checker/maintenance.go b/internal/checker/maintenance.go index 48e1e26..ebbd23a 100644 --- a/internal/checker/maintenance.go +++ b/internal/checker/maintenance.go @@ -3,7 +3,6 @@ package checker import ( "context" "fmt" - "slices" "time" "go.uber.org/zap" @@ -56,16 +55,12 @@ func (st *MntStatusHistory) setStatus(status event.Status) { func (ch *Checker) CheckMaintenance() error { ch.log.Info("check maintenances statuses") - if ch.lastMntID == 0 { - ch.log.Info("no last completed maintenance, starting from the beginning") - } - maintenances, err := ch.db.GetMaintenances(ch.lastMntID) + maintenances, err := ch.db.GetMaintenances() if err != nil { return err } - var activeMaintenances []uint for _, mn := range maintenances { // Draft maintenances are not processed by the checker — they await // manual approval (reviewed) or rejection (cancelled) via the API. @@ -73,37 +68,19 @@ func (ch *Checker) CheckMaintenance() error { continue } - if processErr := ch.processMaintenance(mn, &activeMaintenances); processErr != nil { + if processErr := ch.processMaintenance(mn); processErr != nil { ch.log.Error("failed to process maintenance", zap.Uint("mntID", mn.ID), zap.Error(processErr)) continue } } - if len(activeMaintenances) == 0 { - for _, mn := range maintenances { - if mn.ID > ch.lastMntID { - ch.lastMntID = mn.ID - } - } - ch.log.Debug( - "there are no actual maintenances, set the last ID to the last one", - zap.Uint("lastMntID", ch.lastMntID), - ) - } else { - ch.lastMntID = slices.Min(activeMaintenances) - ch.log.Debug( - "set the last ID to the earliest planned or in_progress maintenance", - zap.Uint("lastMntID", ch.lastMntID), - ) - } - ch.log.Info("finished checking maintenances") return nil } -func (ch *Checker) processMaintenance(mn *db.Incident, activeMaintenances *[]uint) error { +func (ch *Checker) processMaintenance(mn *db.Incident) error { // Refetch immediately before the read-modify-write. The bulk // GetMaintenances above is N items old by the time we reach item N; // using its preloaded state for the version check races concurrent @@ -141,7 +118,6 @@ func (ch *Checker) processMaintenance(mn *db.Incident, activeMaintenances *[]uin ch.notifier.Notify() // wake the worker after the commit } - trackActiveMaintenance(actualStatus, mn.ID, activeMaintenances) return nil } @@ -160,12 +136,6 @@ func (ch *Checker) evaluateAndFixMntStatus(mn *db.Incident) event.Status { return actualStatus } -func trackActiveMaintenance(status event.Status, id uint, activeMaintenances *[]uint) { - if status == event.MaintenancePlanned || status == event.MaintenanceInProgress { - *activeMaintenances = append(*activeMaintenances, id) - } -} - func calculateMntStatusHistory(mn *db.Incident) *MntStatusHistory { sHistory := &MntStatusHistory{} for _, st := range mn.Statuses { diff --git a/internal/db/errors.go b/internal/db/errors.go index 43002bf..6990f98 100644 --- a/internal/db/errors.go +++ b/internal/db/errors.go @@ -12,3 +12,4 @@ var ErrDBIncidentFilterActiveFalse = errors.New("filter for inactive incidents i var ErrVersionConflict = errors.New("version conflict") var ErrNotificationSchemaMissing = errors.New( "notification_outbox table is missing: apply the pending database migrations") +var ErrLockBusy = errors.New("advisory lock is held by another instance") diff --git a/internal/db/event_types.go b/internal/db/event_types.go index ff99a8a..1aaefba 100644 --- a/internal/db/event_types.go +++ b/internal/db/event_types.go @@ -11,7 +11,7 @@ import ( ) // getEventsByType lists events of a single type with their update history. -func (db *DB) getEventsByType(eventType incident.Type, after uint, order entsql.OrderTermOption) ([]*Incident, error) { +func (db *DB) getEventsByType(eventType incident.Type, order entsql.OrderTermOption) ([]*Incident, error) { ctx := context.Background() query := db.e.Incident.Query(). @@ -21,10 +21,6 @@ func (db *DB) getEventsByType(eventType incident.Type, after uint, order entsql. }). Order(incident.ByID(order)) - if after > 0 { - query = query.Where(incident.IDGTE(int(after))) - } - rows, err := query.All(ctx) if err != nil { return nil, err diff --git a/internal/db/info.go b/internal/db/info.go index 33a1aae..d05ab53 100644 --- a/internal/db/info.go +++ b/internal/db/info.go @@ -6,6 +6,6 @@ import ( "github.com/stackmon/otc-status-dashboard/ent/incident" ) -func (db *DB) GetInfoEvents(after uint) ([]*Incident, error) { - return db.getEventsByType(incident.TypeInfo, after, entsql.OrderDesc()) +func (db *DB) GetInfoEvents() ([]*Incident, error) { + return db.getEventsByType(incident.TypeInfo, entsql.OrderDesc()) } diff --git a/internal/db/lock.go b/internal/db/lock.go new file mode 100644 index 0000000..a2a5a65 --- /dev/null +++ b/internal/db/lock.go @@ -0,0 +1,38 @@ +package db + +import ( + "context" +) + +// WithAdvisoryLock runs fn while holding a non-blocking session-level advisory +// lock on key. The lock is acquired and released on the same dedicated +// connection: pg_advisory_lock is session-scoped, so releasing it from a +// different pooled connection would leave the lock held until that connection +// is closed. +// +// The lock is released even when ctx is cancelled, and the connection is +// returned to the pool afterwards. If another session already holds the lock, +// ErrLockBusy is returned without waiting. +func (db *DB) WithAdvisoryLock(ctx context.Context, key int64, fn func(context.Context) error) error { + conn, err := db.sql.Conn(ctx) + if err != nil { + return err + } + defer conn.Close() + + var got bool + if err = conn.QueryRowContext(ctx, "SELECT pg_try_advisory_lock($1)", key).Scan(&got); err != nil { + return err + } + if !got { + return ErrLockBusy + } + + defer func() { + // Unlock must run even after ctx is cancelled, otherwise the lock + // stays held for the lifetime of the connection. + _, _ = conn.ExecContext(context.WithoutCancel(ctx), "SELECT pg_advisory_unlock($1)", key) + }() + + return fn(ctx) +} diff --git a/internal/db/maintenances.go b/internal/db/maintenances.go index b29f6ac..b0ba160 100644 --- a/internal/db/maintenances.go +++ b/internal/db/maintenances.go @@ -6,6 +6,6 @@ import ( "github.com/stackmon/otc-status-dashboard/ent/incident" ) -func (db *DB) GetMaintenances(after uint) ([]*Incident, error) { - return db.getEventsByType(incident.TypeMaintenance, after, entsql.OrderAsc()) +func (db *DB) GetMaintenances() ([]*Incident, error) { + return db.getEventsByType(incident.TypeMaintenance, entsql.OrderAsc()) } diff --git a/tests/advisory_lock_test.go b/tests/advisory_lock_test.go new file mode 100644 index 0000000..8d042f7 --- /dev/null +++ b/tests/advisory_lock_test.go @@ -0,0 +1,102 @@ +package tests + +import ( + "context" + "errors" + "sync" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/stackmon/otc-status-dashboard/internal/db" +) + +func TestWithAdvisoryLock_MutualExclusion(t *testing.T) { + d, _ := newNotifDB(t) + + err := d.WithAdvisoryLock(context.Background(), 9001, func(context.Context) error { + err := d.WithAdvisoryLock(context.Background(), 9001, func(context.Context) error { + return nil + }) + assert.ErrorIs(t, err, db.ErrLockBusy, "same session cannot re-acquire the lock") + return nil + }) + require.NoError(t, err) +} + +func TestWithAdvisoryLock_ConcurrentHolders(t *testing.T) { + d, _ := newNotifDB(t) + + const key = int64(9001) + const workers = 4 + + var mu sync.Mutex + running := 0 + maxRunning := 0 + + run := func() error { + return d.WithAdvisoryLock(context.Background(), key, func(context.Context) error { + mu.Lock() + running++ + if running > maxRunning { + maxRunning = running + } + mu.Unlock() + + time.Sleep(50 * time.Millisecond) + + mu.Lock() + running-- + mu.Unlock() + return nil + }) + } + + var wg sync.WaitGroup + for i := 0; i < workers; i++ { + wg.Add(1) + go func() { + defer wg.Done() + for i := 0; i < 3; i++ { + err := run() + if err != nil && !errors.Is(err, db.ErrLockBusy) { + t.Error(err) + return + } + } + }() + } + wg.Wait() + + assert.Equal(t, 1, maxRunning, "the lock must be held by at most one goroutine at a time") +} + +func TestWithAdvisoryLock_DistinctKeys(t *testing.T) { + d, _ := newNotifDB(t) + + err := d.WithAdvisoryLock(context.Background(), 9001, func(context.Context) error { + return d.WithAdvisoryLock(context.Background(), 9002, func(context.Context) error { + return nil + }) + }) + require.NoError(t, err, "different keys must not block each other") +} + +func TestWithAdvisoryLock_ReleasedAfterCancelledFn(t *testing.T) { + d, _ := newNotifDB(t) + + ctx, cancel := context.WithCancel(context.Background()) + err := d.WithAdvisoryLock(ctx, 9001, func(context.Context) error { + cancel() + return context.Canceled + }) + assert.ErrorIs(t, err, context.Canceled) + + // The lock must have been released despite the cancelled context. + err = d.WithAdvisoryLock(context.Background(), 9001, func(context.Context) error { + return nil + }) + require.NoError(t, err, "lock must be released even when the context was cancelled") +} diff --git a/tests/checker_notifications_test.go b/tests/checker_notifications_test.go index 6c5a62d..fda31ec 100644 --- a/tests/checker_notifications_test.go +++ b/tests/checker_notifications_test.go @@ -11,8 +11,24 @@ import ( "github.com/stackmon/otc-status-dashboard/internal/conf" "github.com/stackmon/otc-status-dashboard/internal/db" "github.com/stackmon/otc-status-dashboard/internal/event" + "github.com/stackmon/otc-status-dashboard/internal/notification" ) +// newTestChecker builds a checker on a fresh pool and a publisher wired to +// the same pool, mirroring the app's wiring. +func newTestChecker(t *testing.T) *checker.Checker { + t.Helper() + + d, err := db.New(&conf.Config{DB: databaseURL}) + require.NoError(t, err) + t.Cleanup(func() { _ = d.Close() }) + + ncfg, err := notification.ConfigFromConf(notifCheckerConfig()) + require.NoError(t, err) + + return checker.New(d, zap.NewNop(), notification.NewPublisher(ncfg, d)) +} + func notifCheckerConfig() *conf.Config { return &conf.Config{ DB: databaseURL, @@ -40,9 +56,7 @@ func TestChecker_ReviewedToPlanned_EnqueuesStatusChangedToCreator(t *testing.T) // No outbox rows yet (publisher was off during API calls). require.Equal(t, int64(0), outboxCount(t, g, eventID)) - chk, err := checker.New(notifCheckerConfig(), zap.NewNop()) - require.NoError(t, err) - t.Cleanup(func() { _ = chk.Close() }) + chk := newTestChecker(t) require.NoError(t, chk.CheckMaintenance()) // reviewed -> planned @@ -62,9 +76,7 @@ func TestChecker_NoTransition_EnqueuesNothing(t *testing.T) { g := openRawDB(t) - chk, err := checker.New(notifCheckerConfig(), zap.NewNop()) - require.NoError(t, err) - t.Cleanup(func() { _ = chk.Close() }) + chk := newTestChecker(t) // Planned with a future start date: the checker computes planned again -> no change. require.NoError(t, chk.CheckMaintenance()) From 8b72aa65591f8e6d987662d74476cad4aa6a3dfa Mon Sep 17 00:00:00 2001 From: Aloento <11802769+Aloento@users.noreply.github.com> Date: Sun, 4 Oct 2026 18:53:04 +0200 Subject: [PATCH 2/2] fix(tests): clear the lint failures in the advisory lock tests The new tests broke the check job: govet flagged the inner loop counter shadowing the outer one, and testifylint required require.ErrorIs for the error assertion. Both loops now use bare for-range so there is no counter to shadow, and the assertion uses require. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- tests/advisory_lock_test.go | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/tests/advisory_lock_test.go b/tests/advisory_lock_test.go index 8d042f7..879cf77 100644 --- a/tests/advisory_lock_test.go +++ b/tests/advisory_lock_test.go @@ -55,11 +55,11 @@ func TestWithAdvisoryLock_ConcurrentHolders(t *testing.T) { } var wg sync.WaitGroup - for i := 0; i < workers; i++ { + for range workers { wg.Add(1) go func() { defer wg.Done() - for i := 0; i < 3; i++ { + for range 3 { err := run() if err != nil && !errors.Is(err, db.ErrLockBusy) { t.Error(err) @@ -92,7 +92,7 @@ func TestWithAdvisoryLock_ReleasedAfterCancelledFn(t *testing.T) { cancel() return context.Canceled }) - assert.ErrorIs(t, err, context.Canceled) + require.ErrorIs(t, err, context.Canceled) // The lock must have been released despite the cancelled context. err = d.WithAdvisoryLock(context.Background(), 9001, func(context.Context) error {