* import with two-stages, wait api and prevnts double run for the same site #231

* lint: missing err check on tmp import file removal

* catch SIQQUIT

* add imprter comments

* add new import apis to spec

* add test for form import

* timeout for import wait api

* lint: uncecked errs
This commit is contained in:
Umputun
2018-12-16 23:21:01 -06:00
committed by GitHub
parent ced78f3332
commit 6bc0d7ee37
6 changed files with 301 additions and 27 deletions
+3 -1
View File
@@ -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.
+1 -1
View File
@@ -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
}
+21
View File
@@ -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)
}
+159 -17
View File
@@ -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
}
+115 -5
View File
@@ -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":"<p>test test #1</p>","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":"<p>test test #2</p>","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":"<p>test test #1</p>","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":"<p>test test #1</p>","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":"<p>test test #1</p>","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":"<p>test test #2</p>","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)
+2 -3
View File
@@ -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,