diff --git a/README.md b/README.md index 478b55b2..e275701f 100644 --- a/README.md +++ b/README.md @@ -579,7 +579,9 @@ Sort can be `time`, `active` or `score`. Supported sort order with prefix -/+, i } ``` * `GET /api/v1/admin/export?site=side-id&mode=[stream|file]` - export all comments to json stream or gz file. -* `POST /api/v1/admin/import?site=side-id` - import comments from the backup. +* `POST /api/v1/admin/import?site=side-id` - import comments from the backup, uses post body. +* `POST /api/v1/admin/import/form?site=side-id` - import comments from the backup, user post form. +* `GET /api/v1/admin/import/wait?site=side-id` - wait for import completeion. * `PUT /api/v1/admin/pin/{id}?site=site-id&url=post-url&pin=1` - pin or unpin comment. * `GET /api/v1/admin/user/{userid}?site=site-id` - get user's info. * `DELETE /api/v1/admin/user/{userid}?site=site-id` - delete all user's comments. diff --git a/backend/app/cmd/import.go b/backend/app/cmd/import.go index d81d4b7c..2e34c445 100644 --- a/backend/app/cmd/import.go +++ b/backend/app/cmd/import.go @@ -62,7 +62,7 @@ func (ic *ImportCommand) Execute(args []string) error { return errors.Wrap(err, "can't get response from importer") } - log.Printf("[INFO] import completed, status=%d, %s", resp.StatusCode, string(body)) + log.Printf("[INFO] completed, status=%d, %s", resp.StatusCode, string(body)) return nil } diff --git a/backend/app/main.go b/backend/app/main.go index 33151e55..92c796cc 100644 --- a/backend/app/main.go +++ b/backend/app/main.go @@ -4,6 +4,9 @@ import ( "fmt" "log" "os" + "os/signal" + "runtime" + "syscall" "github.com/hashicorp/logutils" "github.com/jessevdk/go-flags" @@ -72,3 +75,21 @@ func setupLog(dbg bool) { } log.SetOutput(filter) } + +func init() { + // catch SIGQUIT and print stack traces + sigChan := make(chan os.Signal) + go func() { + for range sigChan { + log.Print("[INFO] SIGQUIT detected") + maxSize := 5 * 1024 * 1024 + stacktrace := make([]byte, maxSize) + length := runtime.Stack(stacktrace, true) + if length > maxSize { + length = maxSize + } + fmt.Println(string(stacktrace[:length])) + } + }() + signal.Notify(sigChan, syscall.SIGQUIT) +} diff --git a/backend/app/rest/api/migrator.go b/backend/app/rest/api/migrator.go index b781cc20..2c489023 100644 --- a/backend/app/rest/api/migrator.go +++ b/backend/app/rest/api/migrator.go @@ -2,14 +2,19 @@ package api import ( "compress/gzip" + "context" "fmt" "io" + "io/ioutil" "log" "net/http" + "os" + "sync" "time" "github.com/go-chi/chi" "github.com/go-chi/render" + "github.com/pkg/errors" "github.com/umputun/remark/backend/app/migrator" "github.com/umputun/remark/backend/app/rest" @@ -24,6 +29,9 @@ type Migrator struct { WordPressImporter migrator.Importer NativeExported migrator.Exporter KeyStore KeyStore + + busy map[string]bool + lock sync.Mutex } // KeyStore defines sub-interface for consumers needed just a key @@ -33,7 +41,10 @@ type KeyStore interface { func (m *Migrator) withRoutes(router chi.Router) chi.Router { router.Get("/export", m.exportCtrl) + router.Post("/import", m.importCtrl) + router.Post("/import/form", m.importFormCtrl) + router.Get("/import/wait", m.importWaitCtrl) return router } @@ -43,26 +54,82 @@ func (m *Migrator) importCtrl(w http.ResponseWriter, r *http.Request) { siteID := r.URL.Query().Get("site") - var importer migrator.Importer - switch r.URL.Query().Get("provider") { - case "disqus": - importer = m.DisqusImporter - case "wordpress": - importer = m.WordPressImporter - default: - importer = m.NativeImporter - } - - log.Printf("[DEBUG] import request for site=%s, provider=%s", siteID, r.URL.Query().Get("provider")) - size, err := importer.Import(r.Body, siteID) - if err != nil { - rest.SendErrorJSON(w, r, http.StatusBadRequest, err, "import failed") + if m.isBusy(siteID) { + rest.SendErrorJSON(w, r, http.StatusConflict, errors.New("already running"), "import rejected") return } - m.Cache.Flush(cache.Flusher(siteID).Scopes(siteID)) - render.Status(r, http.StatusCreated) - render.JSON(w, r, JSON{"status": "ok", "size": size}) + tmpfile, err := m.saveTemp(r.Body) + if err != nil { + rest.SendErrorJSON(w, r, http.StatusInternalServerError, err, "can't save request to temp file") + return + } + + go m.runImport(siteID, r.URL.Query().Get("provider"), tmpfile) // import runs in background and sets busy flag for site + + render.Status(r, http.StatusAccepted) + render.JSON(w, r, JSON{"status": "import request accepted"}) +} + +// POST /import/form?secret=key&site=site-id&provider=disqus|remark|wordpress +// imports comments from form body. +func (m *Migrator) importFormCtrl(w http.ResponseWriter, r *http.Request) { + siteID := r.URL.Query().Get("site") + + if m.isBusy(siteID) { + rest.SendErrorJSON(w, r, http.StatusConflict, errors.New("already running"), "import rejected") + return + } + + if err := r.ParseMultipartForm(20 * 1024 * 1024); err != nil { // 20M max memory, if bigger will make a file + rest.SendErrorJSON(w, r, http.StatusInternalServerError, err, "can't parse multipart form") + return + } + + file, _, err := r.FormFile("file") + if err != nil { + rest.SendErrorJSON(w, r, http.StatusInternalServerError, err, "can't get import from the request") + return + } + defer func() { _ = file.Close() }() + + tmpfile, err := m.saveTemp(file) + if err != nil { + rest.SendErrorJSON(w, r, http.StatusInternalServerError, err, "can't save request to temp file") + return + } + + go m.runImport(siteID, r.URL.Query().Get("provider"), tmpfile) // import runs in background and sets busy flag for site + + render.Status(r, http.StatusAccepted) + render.JSON(w, r, JSON{"status": "import request accepted"}) +} + +func (m *Migrator) importWaitCtrl(w http.ResponseWriter, r *http.Request) { + siteID := r.URL.Query().Get("site") + timeOut := time.Minute * 15 + if v := r.URL.Query().Get("timeout"); v != "" { + if vv, e := time.ParseDuration(v); e == nil { + timeOut = vv + } + } + + ctx, cancel := context.WithTimeout(context.Background(), timeOut) + defer cancel() + for { + if !m.isBusy(siteID) { + break + } + select { + case <-ctx.Done(): + render.Status(r, http.StatusGatewayTimeout) + render.JSON(w, r, JSON{"status": "timeout expired", "site_id": siteID}) + return + case <-time.After(100 * time.Millisecond): + } + } + render.Status(r, http.StatusOK) + render.JSON(w, r, JSON{"status": "completed", "site_id": siteID}) } // GET /export?site=site-id&secret=12345&?mode=file|stream @@ -91,3 +158,78 @@ func (m *Migrator) exportCtrl(w http.ResponseWriter, r *http.Request) { return } } + +// runImport reads from tmpfile and import for given siteID and provider +func (m *Migrator) runImport(siteID string, provider string, tmpfile string) { + m.setBusy(siteID, true) + + defer func() { + m.setBusy(siteID, false) + if err := os.Remove(tmpfile); err != nil { + log.Printf("[WARN] failed to remove tmp file %s, %v", tmpfile, err) + } + }() + + var importer migrator.Importer + switch provider { + case "disqus": + importer = m.DisqusImporter + case "wordpress": + importer = m.WordPressImporter + default: + importer = m.NativeImporter + } + log.Printf("[DEBUG] import request for site=%s, provider=%s", siteID, provider) + + fh, err := os.Open(tmpfile) + if err != nil { + log.Printf("[WARN] import failed, %v", err) + return + } + + size, err := importer.Import(fh, siteID) + if err != nil { + log.Printf("[WARN] import failed, %v", err) + return + } + m.Cache.Flush(cache.Flusher(siteID).Scopes(siteID)) + log.Printf("[DEBUG] import request completed. site=%s, provider=%s, comments=%d", siteID, provider, size) +} + +// saveTemp reads from reader and saves to temp file +func (m *Migrator) saveTemp(r io.Reader) (string, error) { + tmpfile, err := ioutil.TempFile("", "remark42_import") + if err != nil { + return "", errors.Wrap(err, "can't make temp file") + } + + if _, err = io.Copy(tmpfile, r); err != nil { + return "", errors.Wrap(err, "can't copy to temp file") + } + + if err = tmpfile.Close(); err != nil { + return "", errors.Wrap(err, "can't close temp file") + } + + return tmpfile.Name(), nil +} + +// isBusy checks busy flag from the map by siteID as key +func (m *Migrator) isBusy(siteID string) bool { + m.lock.Lock() + defer m.lock.Unlock() + if m.busy == nil { + m.busy = map[string]bool{} + } + return m.busy[siteID] +} + +// setBusy sets/resets busy flag to the map by siteID as key +func (m *Migrator) setBusy(siteID string, status bool) { + m.lock.Lock() + defer m.lock.Unlock() + if m.busy == nil { + m.busy = map[string]bool{} + } + m.busy[siteID] = status +} diff --git a/backend/app/rest/api/migrator_test.go b/backend/app/rest/api/migrator_test.go index c192bd94..031e8d04 100644 --- a/backend/app/rest/api/migrator_test.go +++ b/backend/app/rest/api/migrator_test.go @@ -1,9 +1,13 @@ package api import ( + "bytes" "compress/gzip" "encoding/json" + "fmt" + "io" "io/ioutil" + "mime/multipart" "net/http" "net/http/httptest" "os" @@ -38,13 +42,52 @@ func TestMigrator_Import(t *testing.T) { assert.Nil(t, err) resp, err := client.Do(req) assert.Nil(t, err) - assert.Equal(t, http.StatusCreated, resp.StatusCode) + assert.Equal(t, http.StatusAccepted, resp.StatusCode) b, err := ioutil.ReadAll(resp.Body) assert.Nil(t, err) - assert.Equal(t, `{"size":2,"status":"ok"}`+"\n", string(b)) + assert.Equal(t, "{\"status\":\"import request accepted\"}\n", string(b)) + + client = &http.Client{Timeout: 10 * time.Second} + req, err = http.NewRequest("GET", ts.URL+"/import/wait?site=radio-t", nil) + req.SetBasicAuth("dev", "password") + assert.NoError(t, err) + resp, err = client.Do(req) + assert.Equal(t, 200, resp.StatusCode) } +func TestMigrator_ImportForm(t *testing.T) { + srv, _, ts := prepImportSrv(t) + assert.NotNil(t, srv) + defer cleanupImportSrv(srv, ts) + + r := strings.NewReader(`{"id":"2aa0478c-df1b-46b1-b561-03d507cf482c","pid":"","text":"
test test #1
","user":{"name":"developer one","id":"dev","picture":"/api/v1/avatar/remark.image","profile":"https://remark42.com","admin":true,"ip":"ae12fe3b5f129b5cc4cdd2b136b7b7947c4d2741"},"locator":{"site":"radio-t","url":"https://radio-t.com/blah1"},"score":0,"votes":{},"time":"2018-04-30T01:37:00.849053725-05:00"} + {"id":"83fd97fd-ff64-48d1-9fb7-ca7769c77037","pid":"p1","text":"test test #2
","user":{"name":"developer one","id":"dev","picture":"/api/v1/avatar/remark.image","profile":"https://remark42.com","admin":true,"ip":"ae12fe3b5f129b5cc4cdd2b136b7b7947c4d2741"},"locator":{"site":"radio-t","url":"https://radio-t.com/blah2"},"score":0,"votes":{},"time":"2018-04-30T01:37:00.861387771-05:00"}`) + + bodyBuf := &bytes.Buffer{} + bodyWriter := multipart.NewWriter(bodyBuf) + fileWriter, err := bodyWriter.CreateFormFile("file", "import.json") + require.NoError(t, err) + _, err = io.Copy(fileWriter, r) + require.NoError(t, err) + contentType := bodyWriter.FormDataContentType() + bodyWriter.Close() + + resp, err := http.Post(ts.URL+"/import/form?site=radio-t&provider=native&secret=123456", contentType, bodyBuf) + assert.Nil(t, err) + assert.Equal(t, http.StatusAccepted, resp.StatusCode) + + b, err := ioutil.ReadAll(resp.Body) + assert.Nil(t, err) + assert.Equal(t, "{\"status\":\"import request accepted\"}\n", string(b)) + + client := &http.Client{Timeout: 10 * time.Second} + req, err := http.NewRequest("GET", ts.URL+"/import/wait?site=radio-t", nil) + req.SetBasicAuth("dev", "password") + assert.NoError(t, err) + resp, err = client.Do(req) + assert.Equal(t, 200, resp.StatusCode) +} func TestMigrator_ImportFromWP(t *testing.T) { srv, ds, ts := prepImportSrv(t) assert.NotNil(t, srv) @@ -58,11 +101,18 @@ func TestMigrator_ImportFromWP(t *testing.T) { req.Header.Add("Content-Type", "application/xml; charset=utf-8") resp, err := client.Do(req) assert.Nil(t, err) - assert.Equal(t, http.StatusCreated, resp.StatusCode) + assert.Equal(t, http.StatusAccepted, resp.StatusCode) b, err := ioutil.ReadAll(resp.Body) assert.Nil(t, err) - assert.Equal(t, `{"size":3,"status":"ok"}`+"\n", string(b)) + assert.Equal(t, "{\"status\":\"import request accepted\"}\n", string(b)) + + client = &http.Client{Timeout: 10 * time.Second} + req, err = http.NewRequest("GET", ts.URL+"/import/wait?site=radio-t", nil) + req.SetBasicAuth("dev", "password") + assert.NoError(t, err) + resp, err = client.Do(req) + assert.Equal(t, 200, resp.StatusCode) assert.NoError(t, ds.Interface.Close()) @@ -97,6 +147,59 @@ func TestMigrator_ImportRejected(t *testing.T) { assert.Equal(t, http.StatusUnauthorized, resp.StatusCode) } +func TestMigrator_ImportDouble(t *testing.T) { + srv, _, ts := prepImportSrv(t) + assert.NotNil(t, srv) + defer cleanupImportSrv(srv, ts) + + tmpl := `{"id":"%d","pid":"","text":"test test #1
","user":{"name":"developer one","id":"dev","picture":"/api/v1/avatar/remark.image","profile":"https://remark42.com","admin":true,"ip":"ae12fe3b5f129b5cc4cdd2b136b7b7947c4d2741"},"locator":{"site":"radio-t","url":"https://radio-t.com/blah1"},"score":0,"votes":{},"time":"2018-04-30T01:37:00.849053725-05:00"}` + recs := []string{} + for i := 0; i < 1000; i++ { + recs = append(recs, fmt.Sprintf(tmpl, i)) + } + r := strings.NewReader(strings.Join(recs, "\n")) // reader with 10k records + client := &http.Client{Timeout: 1 * time.Second} + req, err := http.NewRequest("POST", ts.URL+"/import?site=radio-t&provider=native&secret=123456", r) + assert.Nil(t, err) + resp, err := client.Do(req) + assert.Nil(t, err) + assert.Equal(t, http.StatusAccepted, resp.StatusCode) + + client = &http.Client{Timeout: 1 * time.Second} + req, err = http.NewRequest("POST", ts.URL+"/import?site=radio-t&provider=native&secret=123456", r) + assert.Nil(t, err) + resp, err = client.Do(req) + assert.Nil(t, err) + assert.Equal(t, http.StatusConflict, resp.StatusCode) + +} + +func TestMigrator_ImportWaitExpired(t *testing.T) { + srv, _, ts := prepImportSrv(t) + assert.NotNil(t, srv) + defer cleanupImportSrv(srv, ts) + + tmpl := `{"id":"%d","pid":"","text":"test test #1
","user":{"name":"developer one","id":"dev","picture":"/api/v1/avatar/remark.image","profile":"https://remark42.com","admin":true,"ip":"ae12fe3b5f129b5cc4cdd2b136b7b7947c4d2741"},"locator":{"site":"radio-t","url":"https://radio-t.com/blah1"},"score":0,"votes":{},"time":"2018-04-30T01:37:00.849053725-05:00"}` + recs := []string{} + for i := 0; i < 1000; i++ { + recs = append(recs, fmt.Sprintf(tmpl, i)) + } + r := strings.NewReader(strings.Join(recs, "\n")) // reader with 10k records + client := &http.Client{Timeout: 1 * time.Second} + req, err := http.NewRequest("POST", ts.URL+"/import?site=radio-t&provider=native&secret=123456", r) + require.Nil(t, err) + resp, err := client.Do(req) + assert.Nil(t, err) + assert.Equal(t, http.StatusAccepted, resp.StatusCode) + + client = &http.Client{Timeout: 10 * time.Second} + req, err = http.NewRequest("GET", ts.URL+"/import/wait?site=radio-t&timeout=100ms", nil) + req.SetBasicAuth("dev", "password") + assert.NoError(t, err) + resp, err = client.Do(req) + assert.Equal(t, http.StatusGatewayTimeout, resp.StatusCode) +} + func TestMigrator_Export(t *testing.T) { srv, _, ts := prepImportSrv(t) assert.NotNil(t, srv) @@ -105,12 +208,19 @@ func TestMigrator_Export(t *testing.T) { r := strings.NewReader(`{"id":"2aa0478c-df1b-46b1-b561-03d507cf482c","pid":"","text":"test test #1
","user":{"name":"developer one","id":"dev","picture":"/api/v1/avatar/remark.image","profile":"https://remark42.com","admin":true,"ip":"ae12fe3b5f129b5cc4cdd2b136b7b7947c4d2741"},"locator":{"site":"radio-t","url":"https://radio-t.com/blah1"},"score":0,"votes":{},"time":"2018-04-30T01:37:00.849053725-05:00"} {"id":"83fd97fd-ff64-48d1-9fb7-ca7769c77037","pid":"p1","text":"test test #2
","user":{"name":"developer one","id":"dev","picture":"/api/v1/avatar/remark.image","profile":"https://remark42.com","admin":true,"ip":"ae12fe3b5f129b5cc4cdd2b136b7b7947c4d2741"},"locator":{"site":"radio-t","url":"https://radio-t.com/blah2"},"score":0,"votes":{},"time":"2018-04-30T01:37:00.861387771-05:00"}`) + // import comments first client := &http.Client{Timeout: 1 * time.Second} req, err := http.NewRequest("POST", ts.URL+"/import?site=radio-t&provider=native&secret=123456", r) require.Nil(t, err) resp, err := client.Do(req) require.Nil(t, err) - require.Equal(t, http.StatusCreated, resp.StatusCode) + require.Equal(t, http.StatusAccepted, resp.StatusCode) + client = &http.Client{Timeout: 10 * time.Second} + req, err = http.NewRequest("GET", ts.URL+"/import/wait?site=radio-t", nil) + req.SetBasicAuth("dev", "password") + assert.NoError(t, err) + resp, err = client.Do(req) + assert.Equal(t, 200, resp.StatusCode) // check file mode req, err = http.NewRequest("GET", ts.URL+"/export?mode=file&site=radio-t&secret=123456", nil) diff --git a/backend/app/rest/api/ssl.go b/backend/app/rest/api/ssl.go index 40cd5664..a087a613 100644 --- a/backend/app/rest/api/ssl.go +++ b/backend/app/rest/api/ssl.go @@ -113,11 +113,10 @@ func makeTLSConfig() *tls.Config { CipherSuites: []uint16{ tls.TLS_ECDHE_ECDSA_WITH_AES_256_GCM_SHA384, tls.TLS_ECDHE_RSA_WITH_AES_256_GCM_SHA384, - tls.TLS_ECDHE_ECDSA_WITH_CHACHA20_POLY1305, - tls.TLS_ECDHE_RSA_WITH_CHACHA20_POLY1305, + // tls.TLS_ECDHE_ECDSA_WITH_CHACHA20_POLY1305, + // tls.TLS_ECDHE_RSA_WITH_CHACHA20_POLY1305, tls.TLS_ECDHE_ECDSA_WITH_AES_128_GCM_SHA256, tls.TLS_ECDHE_RSA_WITH_AES_128_GCM_SHA256, - tls.TLS_ECDHE_ECDSA_WITH_AES_256_CBC_SHA, }, MinVersion: tls.VersionTLS12,