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
114 changes: 97 additions & 17 deletions cmd/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -14,10 +14,20 @@ import (
"github.com/stackmon/otc-status-dashboard/internal/app"
"github.com/stackmon/otc-status-dashboard/internal/checker"
"github.com/stackmon/otc-status-dashboard/internal/conf"
"github.com/stackmon/otc-status-dashboard/internal/scheduler"
)

// shutdownTimeout bounds the in-flight request drain after SIGTERM.
const shutdownTimeout = 15 * time.Second
const (
// shutdownTimeout bounds the in-flight request drain after SIGTERM.
shutdownTimeout = 15 * time.Second
// taskStopTimeout bounds the wait for in-flight scheduled tasks and the
// notification worker during shutdown.
taskStopTimeout = 30 * time.Second

scanInterval = time.Minute * 2
sweepInterval = time.Minute * 5
retentionInterval = time.Hour * 24
)

func main() {
c, err := conf.LoadConf()
Expand All @@ -33,35 +43,105 @@ func main() {
logger.Fatal("fail to init app", zap.Error(err))
}

ch := checker.New(s.DB, logger, s.Publisher())

ctx, done := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
defer done()

go func() {
if err = s.Run(); err != nil && !errors.Is(err, http.ErrServerClosed) {
logger.Fatal("app is failed to run", zap.Error(err))
}
}()

go ch.Run()
sched := newScheduler(s, logger)
go runServer(s, logger)
sched.Run(ctx)
workerDone := startWorker(ctx, s)

<-ctx.Done()
s.Log.Info("shutdown app")
shutdown(s, sched, logger, workerDone)
logger.Info("app exited")
}

// newScheduler registers every periodic task. The advisory lock that keeps a task
// single-replica is taken by the scheduler, not by the task body.
func newScheduler(s *app.App, logger *zap.Logger) *scheduler.Scheduler {
ch := checker.New(s.DB, logger, s.Publisher())

// Stop the checker before the pool is closed: Check runs synchronously, so
// this waits for an in-flight scan to finish before App.Shutdown closes the
// database pool.
ch.Shutdown()
sched := scheduler.New(s.DB, logger)
sched.Register("scan", scanInterval, scheduler.KeyScan, func(ctx context.Context) error {
if err := ch.Check(ctx); err != nil {
return err
}
s.Publisher().Notify()
return nil
})
if w := s.Worker(); w != nil {
sched.Register("notify_sweep", sweepInterval, scheduler.KeyNotifySweep, w.Drain)
sched.Register("retention", retentionInterval, scheduler.KeyRetention, w.RunRetention)
}
return sched
}

func runServer(s *app.App, logger *zap.Logger) {
if err := s.Run(); err != nil && !errors.Is(err, http.ErrServerClosed) {
logger.Fatal("app is failed to run", zap.Error(err))
}
}

// startWorker runs the notification worker on ctx and returns a channel closed
// when it has stopped, or nil when notifications are disabled.
func startWorker(ctx context.Context, s *app.App) chan struct{} {
w := s.Worker()
if w == nil {
return nil
}
done := make(chan struct{})
go func() {
w.Run(ctx)
close(done)
}()
return done
}

func shutdown(s *app.App, sched *scheduler.Scheduler, logger *zap.Logger, workerDone chan struct{}) {
// The signal context is already cancelled, so the shutdown needs its own
// deadline to drain in-flight requests.
shutdownCtx, cancel := context.WithTimeout(context.Background(), shutdownTimeout)
defer cancel()

if err = s.Shutdown(shutdownCtx); err != nil {
// A scan round holds a dedicated connection, so the pool must outlive the
// scheduled work. The scheduler and the worker get independent deadlines so
// a slow scheduler stop cannot starve the worker wait.
schedErr := stopScheduler(sched)
workerErr := waitWorker(workerDone, logger)

if err := s.Shutdown(shutdownCtx); err != nil {
logger.Error("app shutdown failed", zap.Error(err))
}

logger.Info("app exited")
// Never close the pool while in-flight work may still be running.
if schedErr != nil || workerErr != nil {
logger.Error("in-flight work did not stop before the deadline; leaving the database pool open",
zap.Error(errors.Join(schedErr, workerErr)))
return
}
if err := s.DB.Close(); err != nil {
logger.Error("database close failed", zap.Error(err))
}
}

func stopScheduler(sched *scheduler.Scheduler) error {
ctx, cancel := context.WithTimeout(context.Background(), taskStopTimeout)
defer cancel()
return sched.Stop(ctx)
}

func waitWorker(workerDone chan struct{}, logger *zap.Logger) error {
if workerDone == nil {
return nil
}
ctx, cancel := context.WithTimeout(context.Background(), taskStopTimeout)
defer cancel()
select {
case <-workerDone:
return nil
case <-ctx.Done():
logger.Warn("timed out waiting for the notification worker to stop")
return ctx.Err()
}
}
28 changes: 13 additions & 15 deletions internal/app/app.go
Original file line number Diff line number Diff line change
Expand Up @@ -36,9 +36,9 @@ type App struct {
srv *http.Server
// metrics server, listening on its own port (nil when notifications are disabled)
metricsSrv *http.Server
// notification delivery worker (nil when notifications are disabled)
worker *notification.Worker
workerCancel context.CancelFunc
// notification delivery worker (nil when notifications are disabled); its
// lifecycle is owned by main, not by the app
worker *notification.Worker
}

func New(c *conf.Config, log *zap.Logger) (*App, error) {
Expand Down Expand Up @@ -128,18 +128,20 @@ func (a *App) NotifyFunc() func() {
return a.worker.Notify
}

// Worker returns the notification delivery worker, or nil when notifications are
// disabled. main owns its lifecycle: start it before the scheduler and stop it
// before the database pool is closed.
func (a *App) Worker() *notification.Worker {
return a.worker
}

// 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
ctx, a.workerCancel = context.WithCancel(context.Background())
go a.worker.Run(ctx)
}
if a.metricsSrv != nil {
go func() {
a.Log.Info("metrics server started", zap.String("addr", a.metricsSrv.Addr))
Expand All @@ -151,17 +153,13 @@ func (a *App) Run() error {
return a.srv.ListenAndServe()
}

// Shutdown stops the HTTP and metrics listeners. The database pool is owned by
// the caller: main closes it last, once in-flight work has stopped.
func (a *App) Shutdown(ctx context.Context) error {
if a.workerCancel != nil {
a.workerCancel()
}
if a.metricsSrv != nil {
if err := a.metricsSrv.Shutdown(ctx); err != nil {
a.Log.Error("metrics server shutdown", zap.Error(err))
}
}
if err := a.srv.Shutdown(ctx); err != nil {
return err
}
return a.DB.Close()
return a.srv.Shutdown(ctx)
}
91 changes: 21 additions & 70 deletions internal/checker/checker.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,109 +4,60 @@ import (
"context"
"errors"
"sync"
"time"

"go.uber.org/zap"

"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
mu sync.Mutex
cancel context.CancelFunc
done chan struct{}
}

// 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{})}
return &Checker{db: database, log: log, notifier: notifier}
}

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
// Check runs one full scan and returns the combined error of its two halves. It
// is the body of the scheduler's scan task, which holds the advisory lock for the
// whole round. Cancellation is observed only before the round starts: the two
// scans do not take a context yet, so a caller must not close the pool while
// Check is running.
func (ch *Checker) Check(ctx context.Context) error {
if err := ctx.Err(); err != nil {
return err
}
if err != nil {
ch.log.Error("failed to acquire the scan lock", zap.Error(err))
}
}

func (ch *Checker) runScan() {
var wg sync.WaitGroup
var (
wg sync.WaitGroup
mntErr error
infoErr error
)

wg.Add(1)
go func() {
err := ch.CheckMaintenance()
if err != nil {
defer wg.Done()
if err := ch.CheckMaintenance(); err != nil {
ch.log.Error("error to check maintenances", zap.Error(err))
mntErr = err
}
wg.Done()
}()

wg.Add(1)
go func() {
err := ch.CheckInfoEvents()
if err != nil {
defer wg.Done()
if err := ch.CheckInfoEvents(); err != nil {
ch.log.Error("error to check info events", zap.Error(err))
infoErr = err
}
wg.Done()
}()

wg.Wait()
}

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 <-ctx.Done():
close(ch.done)
return
case <-ticker.C:
ch.Check()
}
}
}

// 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")
ch.mu.Lock()
cancel := ch.cancel
ch.mu.Unlock()
if cancel == nil {
return
}
cancel()
<-ch.done
return errors.Join(mntErr, infoErr)
}
Loading
Loading