main split to setup and run stages, with graceful shutdown

This commit is contained in:
Umputun
2018-05-22 20:44:02 -05:00
parent 781bfca637
commit bb4a54ff71
4 changed files with 104 additions and 34 deletions
+67 -25
View File
@@ -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)
}
}
+15 -7
View File
@@ -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))
}
}
+13 -2
View File
@@ -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)
+9
View File
@@ -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)