Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
25 changes: 3 additions & 22 deletions cmd/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand All @@ -55,9 +40,7 @@ func main() {
}
}()

go func() {
ch.Run(stopCh)
}()
go ch.Run()

<-ctx.Done()
s.Log.Info("shutdown app")
Expand All @@ -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")
}
8 changes: 7 additions & 1 deletion internal/app/app.go
Original file line number Diff line number Diff line change
Expand Up @@ -120,14 +120,20 @@ 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
}
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
Expand Down
78 changes: 49 additions & 29 deletions internal/checker/checker.go
Original file line number Diff line number Diff line change
@@ -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)
Expand All @@ -58,35 +75,38 @@ 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()
}
}
}

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
}
27 changes: 1 addition & 26 deletions internal/checker/info.go
Original file line number Diff line number Diff line change
@@ -1,7 +1,6 @@
package checker

import (
"slices"
"time"

"go.uber.org/zap"
Expand Down Expand Up @@ -48,27 +47,21 @@ 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)

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:
Expand All @@ -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
Expand Down
36 changes: 3 additions & 33 deletions internal/checker/maintenance.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,6 @@ package checker
import (
"context"
"fmt"
"slices"
"time"

"go.uber.org/zap"
Expand Down Expand Up @@ -56,54 +55,32 @@ 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.
if mn.Status == event.MaintenancePendingReview {
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
Expand Down Expand Up @@ -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
}

Expand All @@ -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 {
Expand Down
1 change: 1 addition & 0 deletions internal/db/errors.go
Original file line number Diff line number Diff line change
Expand Up @@ -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")
6 changes: 1 addition & 5 deletions internal/db/event_types.go
Original file line number Diff line number Diff line change
Expand Up @@ -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().
Expand All @@ -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
Expand Down
4 changes: 2 additions & 2 deletions internal/db/info.go
Original file line number Diff line number Diff line change
Expand Up @@ -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())
}
Loading
Loading