diff --git a/app/main.go b/app/main.go index 3b86d3d9..152cbe8b 100644 --- a/app/main.go +++ b/app/main.go @@ -1,11 +1,14 @@ package main import ( + "context" "fmt" "log" "net/http" "os" + "os/signal" "strings" + "syscall" "time" "github.com/coreos/bbolt" @@ -22,7 +25,7 @@ import ( "github.com/umputun/remark/app/rest/proxy" ) -var opts struct { +type Opts struct { BoltPath string `long:"bolt" env:"BOLTDB_PATH" default:"./var" description:"parent dir for bolt files"` Sites []string `long:"site" env:"SITE" default:"remark" description:"site names" env-delim:","` RemarkURL string `long:"url" env:"REMARK_URL" default:"https://remark42.com" description:"url to remark"` @@ -56,49 +59,80 @@ var opts struct { WebRoot string `long:"web-root" env:"REMARK_WEB_ROOT" default:"./web" description:"web root directory"` } +var opts Opts var revision = "unknown" +type application struct { + Opts + srv api.Rest + importer api.Import + exporter migrator.Exporter +} + func main() { fmt.Printf("remark %s\n", revision) p := flags.NewParser(&opts, flags.Default) if _, e := p.ParseArgs(os.Args[1:]); e != nil { os.Exit(1) } - - setupLog(opts.Dbg) log.Print("[INFO] started remark") - if err := makeDirs(opts.BoltPath, opts.BackupLocation, opts.AvatarStore); err != nil { - log.Fatalf("[ERROR] can't create directories, %+v", err) + ctx, cancel := context.WithCancel(context.Background()) + + // graceful termination + go func() { + stop := make(chan os.Signal, 1) + signal.Notify(stop, os.Interrupt, syscall.SIGTERM) + <-stop + log.Print("[WARN] interrupt signal") + cancel() + }() + + app, err := Setup(opts) + if err != nil { + log.Fatalf("[ERROR] failed to setup application, %+v", err) } + log.Printf("[INFO] remark terminated %s", Run(ctx, app)) +} - dataStore := makeBoltStore(opts.Sites) - - if opts.DevPasswd != "" { +// Run all application objects +func Run(ctx context.Context, a *application) error { + if a.DevPasswd != "" { log.Printf("[WARN] running in dev mode") } + activateBackup(ctx, a.exporter) // runs in goroutine for each site + go a.importer.Run(a.Port + 1) + go a.srv.Run(opts.Port) + + // shutdown on context cancellation + <-ctx.Done() + a.srv.Shutdown() + a.importer.Shutdown() + + return ctx.Err() +} + +// Setup prepares application and return all active parts +// doesn't start anything +func Setup(opts Opts) (*application, error) { + setupLog(opts.Dbg) + + if err := makeDirs(opts.BoltPath, opts.BackupLocation, opts.AvatarStore); err != nil { + return nil, err + } + dataService := service.DataStore{ - Interface: dataStore, + Interface: makeBoltStore(opts.Sites), EditDuration: 5 * time.Minute, Secret: opts.SecretKey, MaxCommentSize: opts.MaxCommentSize, } - exporter := migrator.Remark{DataStore: &dataService} cache := rest.NewLoadingCache(rest.MaxValueSize(opts.MaxCachedValue), rest.MaxKeys(opts.MaxCachedItems), rest.PostFlushFn(postFlushFn)) - activateBackup(&exporter) - - importSrv := api.Import{ - Version: revision, - Cache: cache, - NativeImporter: &migrator.Remark{DataStore: &dataService}, - DisqusImporter: &migrator.Disqus{DataStore: &dataService}, - SecretKey: opts.SecretKey, - } - go importSrv.Run(opts.Port + 1) + jwtService := auth.NewJWT(opts.SecretKey, strings.HasPrefix(opts.RemarkURL, "https://"), 7*24*time.Hour) avatarProxy := &proxy.Avatar{ StorePath: opts.AvatarStore, @@ -106,12 +140,20 @@ func main() { RemarkURL: strings.TrimSuffix(opts.RemarkURL, "/"), } - jwtService := auth.NewJWT(opts.SecretKey, strings.HasPrefix(opts.RemarkURL, "https://"), 7*24*time.Hour) + exporter := &migrator.Remark{DataStore: &dataService} + + importer := api.Import{ + Version: revision, + Cache: cache, + NativeImporter: &migrator.Remark{DataStore: &dataService}, + DisqusImporter: &migrator.Disqus{DataStore: &dataService}, + SecretKey: opts.SecretKey, + } srv := api.Rest{ Version: revision, DataService: dataService, - Exporter: &exporter, + Exporter: exporter, WebRoot: opts.WebRoot, ImageProxy: proxy.Image{Enabled: opts.ImageProxy, RoutePath: "/api/v1/img", RemarkURL: opts.RemarkURL}, Authenticator: auth.Authenticator{ @@ -124,11 +166,11 @@ func main() { Cache: cache, } srv.ScoreThresholds.Low, srv.ScoreThresholds.Critical = opts.LowScore, opts.CriticalScore - srv.Run(opts.Port) + return &application{srv: srv, importer: importer, exporter: exporter, Opts: opts}, nil } // activateBackup runs background backups for each site -func activateBackup(exporter migrator.Exporter) { +func activateBackup(ctx context.Context, exporter migrator.Exporter) { for _, siteID := range opts.Sites { backup := migrator.AutoBackup{ Exporter: exporter, @@ -137,7 +179,7 @@ func activateBackup(exporter migrator.Exporter) { KeepMax: opts.MaxBackupFiles, Duration: 24 * time.Hour, } - go backup.Do() + go backup.Do(ctx) } } diff --git a/app/migrator/backup.go b/app/migrator/backup.go index 155ef75c..c0fcff59 100644 --- a/app/migrator/backup.go +++ b/app/migrator/backup.go @@ -2,6 +2,7 @@ package migrator import ( "compress/gzip" + "context" "fmt" "io/ioutil" "log" @@ -23,17 +24,24 @@ type AutoBackup struct { } // Do runs daily export to local files, keeps up to keepMax backups for given siteID -func (ab AutoBackup) Do() { +func (ab AutoBackup) Do(ctx context.Context) { log.Printf("[INFO] activate auto-backup for %s", ab.BackupLocation) tick := time.NewTicker(ab.Duration) log.Printf("[DEBUG] first backup at %s", time.Now().Add(ab.Duration)) - for range tick.C { - if _, err := ab.makeBackup(); err != nil { - log.Printf("[WARN] auto-backup for %s failed, %s", ab.SiteID, err) - continue + + for { + select { + case <-tick.C: + if _, err := ab.makeBackup(); err != nil { + log.Printf("[WARN] auto-backup for %s failed, %s", ab.SiteID, err) + continue + } + ab.removeOldBackupFiles() + log.Printf("[DEBUG] next backup at %s", time.Now().Add(ab.Duration)) + case <-ctx.Done(): + log.Printf("[WARN] terminated autobackup for %s", ab.SiteID) + return } - ab.removeOldBackupFiles() - log.Printf("[DEBUG] next backup at %s", time.Now().Add(ab.Duration)) } } diff --git a/app/rest/api/import.go b/app/rest/api/import.go index 2f34da6d..a69d18f9 100644 --- a/app/rest/api/import.go +++ b/app/rest/api/import.go @@ -1,6 +1,7 @@ package api import ( + "context" "fmt" "log" "net/http" @@ -24,6 +25,8 @@ type Import struct { NativeImporter migrator.Importer DisqusImporter migrator.Importer SecretKey string + + httpServer *http.Server } // Run the listener and request's router, activate rest server @@ -31,11 +34,19 @@ type Import struct { func (s *Import) Run(port int) { log.Printf("[INFO] activate import server on port %d", port) router := s.routes() - httpServer := &http.Server{Addr: fmt.Sprintf("127.0.0.1:%d", port), Handler: router} - err := httpServer.ListenAndServe() + s.httpServer = &http.Server{Addr: fmt.Sprintf("127.0.0.1:%d", port), Handler: router} + err := s.httpServer.ListenAndServe() log.Printf("[WARN] http server terminated, %s", err) } +func (s *Import) Shutdown() { + log.Print("[WARN] shutdown import server") + ctx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + s.httpServer.Shutdown(ctx) + log.Print("[DEBUG] shutdown import server completed") +} + func (s Import) routes() chi.Router { router := chi.NewRouter() router.Use(middleware.RealIP, Recoverer) diff --git a/app/rest/api/rest.go b/app/rest/api/rest.go index 717fec68..a36f9e8c 100644 --- a/app/rest/api/rest.go +++ b/app/rest/api/rest.go @@ -2,6 +2,7 @@ package api import ( "bytes" + "context" "crypto/sha1" "encoding/base64" "encoding/json" @@ -74,6 +75,14 @@ func (s *Rest) Run(port int) { log.Printf("[WARN] http server terminated, %s", err) } +func (s *Rest) Shutdown() { + log.Print("[WARN] shutdown rest server") + ctx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + s.httpServer.Shutdown(ctx) + log.Print("[DEBUG] shutdown rest server completed") +} + func (s *Rest) routes() chi.Router { router := chi.NewRouter() router.Use(middleware.RealIP, Recoverer)