simplify import & export

This commit is contained in:
Umputun
2018-12-23 02:23:54 -06:00
parent 094f4c66a1
commit be3d1bdf3e
4 changed files with 121 additions and 105 deletions
+1
View File
@@ -31,6 +31,7 @@ type Store interface {
List(siteID string, limit int, skip int) ([]store.PostInfo, error)
DeleteAll(siteID string) error
Metas(siteID string) (umetas []service.UserMetaData, pmetas []service.PostMetaData, err error)
SetMetas(siteID string, umetas []service.UserMetaData, pmetas []service.PostMetaData) error
}
// ImportParams defines everything needed to run import
+47 -62
View File
@@ -1,10 +1,8 @@
package migrator
import (
"bufio"
"bytes"
"encoding/json"
"fmt"
"io"
"log"
@@ -14,26 +12,28 @@ import (
"github.com/umputun/remark/backend/app/store/service"
)
const (
header = `{"version":1, "comments":[`
metaHeader = "],\n\"meta\":"
footer = `}`
)
// Native implements exporter and importer for internal store format
// {"version": 1, comments:[{...}\n,{}], meta: {meta}}
// each comments starts from the new line
type Native struct {
DataStore Store
}
// Export all comments to writer as json strings. Each comment is one string, separated by "\n"
func (r *Native) Export(w io.Writer, siteID string) (size int, err error) {
type meta struct {
Version int `json:"version"`
Users []service.UserMetaData `json:"users"`
Posts []service.PostMetaData `json:"posts"`
}
if _, err = fmt.Fprintf(w, "%s\n", header); err != nil {
return 0, err
// Export all comments to writer as json strings. Each comment is one string, separated by "\n"
// The final file is a valid json
func (n *Native) Export(w io.Writer, siteID string) (size int, err error) {
if err = n.exportMeta(siteID, w); err != nil {
return 0, errors.Wrapf(err, "failed to export meta for site %s", siteID)
}
topics, err := r.DataStore.List(siteID, 0, 0)
topics, err := n.DataStore.List(siteID, 0, 0)
if err != nil {
return 0, err
}
@@ -42,12 +42,12 @@ func (r *Native) Export(w io.Writer, siteID string) (size int, err error) {
commentsCount := 0
for i := len(topics) - 1; i >= 0; i-- { // topics from List sorted in opposite direction
topic := topics[i]
comments, e := r.DataStore.Find(store.Locator{SiteID: siteID, URL: topic.URL}, "time")
comments, e := n.DataStore.Find(store.Locator{SiteID: siteID, URL: topic.URL}, "time")
if err != nil {
return commentsCount, e
}
for n, comment := range comments {
for _, comment := range comments {
buf := &bytes.Buffer{}
enc := json.NewEncoder(buf)
@@ -56,77 +56,61 @@ func (r *Native) Export(w io.Writer, siteID string) (size int, err error) {
if err = enc.Encode(comment); err != nil {
return commentsCount, errors.Wrapf(err, "can't marshal %v", comments)
}
data := buf.Bytes()
data = bytes.TrimSuffix(data, []byte("\n"))
if _, err = w.Write(data); err != nil {
if _, err = w.Write(buf.Bytes()); err != nil {
return commentsCount, errors.Wrap(err, "can't write comment data")
}
if n < len(comments)-1 || i != 0 { // don't add , on last comment
if _, err = w.Write([]byte(",")); err != nil {
return commentsCount, errors.Wrap(err, "can't write comment separator")
}
}
if _, err = w.Write([]byte("\n")); err != nil {
return commentsCount, errors.Wrap(err, "can't write comment eol")
}
commentsCount++
}
}
log.Printf("[DEBUG] exported %d comments", commentsCount)
err = r.exportMeta(siteID, w)
return commentsCount, err
return commentsCount, nil
}
func (r *Native) exportMeta(siteID string, w io.Writer) (err error) {
if _, err = fmt.Fprintf(w, "%s", metaHeader); err != nil {
return errors.Wrap(err, "can't write meta header")
}
meta := struct {
Users []service.UserMetaData `json:"users"`
Posts []service.PostMetaData `json:"posts"`
}{}
meta.Users, meta.Posts, err = r.DataStore.Metas(siteID)
// exportMeta appends user and post metas to exported stream
func (n *Native) exportMeta(siteID string, w io.Writer) (err error) {
m := meta{Version: 1}
m.Users, m.Posts, err = n.DataStore.Metas(siteID)
if err != nil {
return errors.Wrap(err, "can't get meta")
}
if err := json.NewEncoder(w).Encode(meta); err != nil {
if err := json.NewEncoder(w).Encode(m); err != nil {
return errors.Wrap(err, "can't encode meta")
}
if _, err := fmt.Fprintf(w, "%s\n", footer); err != nil {
return errors.Wrap(err, "can't write footer")
}
return nil
}
// Import comments from json strings produced by Remark.Export
func (r *Native) Import(reader io.Reader, siteID string) (size int, err error) {
func (n *Native) Import(reader io.Reader, siteID string) (size int, err error) {
if err := r.DataStore.DeleteAll(siteID); err != nil {
m := meta{}
dec := json.NewDecoder(reader)
if err = dec.Decode(&m); err != nil {
return 0, errors.Wrapf(err, "failed to import meta for site %s", siteID)
}
if err := n.DataStore.DeleteAll(siteID); err != nil {
return 0, err
}
failed := 0
total, comments := 0, 0
scanner := bufio.NewScanner(reader)
for scanner.Scan() {
rec := scanner.Bytes()
if len(rec) < 3 {
continue
}
total++
for {
comment := store.Comment{}
if err := json.Unmarshal(rec, &comment); err != nil {
err = dec.Decode(&comment)
if err == io.EOF {
break
}
total++
if err != nil {
failed++
log.Printf("[WARN] unmarshal failed for %s, %s", string(rec), err)
continue
}
if _, err := r.DataStore.Create(comment); err != nil {
if _, err := n.DataStore.Create(comment); err != nil {
failed++
log.Printf("[WARN] can't write %+v to store, %s", comment, err)
continue
@@ -136,12 +120,13 @@ func (r *Native) Import(reader io.Reader, siteID string) (size int, err error) {
log.Printf("[DEBUG] imported %d comments", comments)
}
}
if scanner.Err() != nil {
return comments, errors.Wrap(scanner.Err(), "error in scan")
}
if failed > 0 {
return comments, errors.Errorf("failed to save %d comments", failed)
}
log.Printf("[INFO] imported %d comments from %d records", comments, total)
return comments, nil
err = n.DataStore.SetMetas(siteID, m.Users, m.Posts)
return comments, err
}
+47 -43
View File
@@ -6,11 +6,13 @@ import (
"fmt"
"log"
"os"
"strings"
"testing"
"time"
bolt "github.com/coreos/bbolt"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"github.com/umputun/remark/backend/app/store"
"github.com/umputun/remark/backend/app/store/admin"
@@ -36,65 +38,66 @@ func TestNative_Export(t *testing.T) {
c1 := buf.String()
log.Print(c1)
res := struct {
Version int `json:"version"`
Comments []store.Comment `json:"comments"`
Meta struct {
Users []service.UserMetaData `json:"users"`
Posts []service.PostMetaData `json:"posts"`
} `json:"meta"`
dec := json.NewDecoder(strings.NewReader(c1))
meta := struct {
Version int `json:"version"`
Users []service.UserMetaData `json:"users"`
Posts []service.PostMetaData `json:"posts"`
}{}
err = json.Unmarshal([]byte(c1), &res)
assert.NoError(t, err)
assert.Equal(t, 2, len(res.Comments))
assert.Equal(t, "some text, <a href=\"http://radio-t.com\" rel=\"nofollow\">link</a>", res.Comments[0].Text)
require.NoError(t, dec.Decode(&meta), "decode meta")
assert.Equal(t, 2, len(res.Meta.Users))
assert.Equal(t, "user1", res.Meta.Users[0].ID)
assert.Equal(t, false, res.Meta.Users[0].Blocked.Status)
assert.Equal(t, true, res.Meta.Users[0].Verified)
assert.Equal(t, "user2", res.Meta.Users[1].ID)
assert.Equal(t, true, res.Meta.Users[1].Blocked.Status)
assert.Equal(t, false, res.Meta.Users[1].Verified)
assert.Equal(t, 2, len(meta.Users))
assert.Equal(t, "user1", meta.Users[0].ID)
assert.Equal(t, false, meta.Users[0].Blocked.Status)
assert.Equal(t, true, meta.Users[0].Verified)
assert.Equal(t, "user2", meta.Users[1].ID)
assert.Equal(t, true, meta.Users[1].Blocked.Status)
assert.Equal(t, false, meta.Users[1].Verified)
assert.Equal(t, 1, len(res.Meta.Posts))
assert.Equal(t, "https://radio-t.com", res.Meta.Posts[0].URL)
assert.Equal(t, true, res.Meta.Posts[0].ReadOnly)
assert.Equal(t, 1, len(meta.Posts))
assert.Equal(t, "https://radio-t.com", meta.Posts[0].URL)
assert.Equal(t, true, meta.Posts[0].ReadOnly)
comments := [3]store.Comment{}
assert.NoError(t, dec.Decode(&comments[0]), "decode comment 0")
assert.NoError(t, dec.Decode(&comments[1]), "decode comment 0")
assert.Error(t, dec.Decode(&comments[2]), "EOF")
assert.Equal(t, "some text, <a href=\"http://radio-t.com\" rel=\"nofollow\">link</a>", comments[0].Text)
}
func TestNative_Import(t *testing.T) {
defer os.Remove(testDb)
r1 := `{"id":"efbc17f177ee1a1c0ee6e1e025749966ec071adc","pid":"","text":"some text, <a href=\"http://radio-t.com\" rel=\"nofollow\">link</a>","user":{"name":"user name","id":"user1","picture":"","profile":"","admin":false},"locator":{"site":"radio-t","url":"https://radio-t.com"},"score":0,"votes":{},"time":"2017-12-20T15:18:22-06:00"}` + "\n"
r2 := `{"id":"afbc17f177ee1a1c0ee6e1e025749966ec071adc","pid":"efbc17f177ee1a1c0ee6e1e025749966ec071adc","text":"some text2, <a href=\"http://radio-t.com\" rel=\"nofollow\">link</a>","user":{"name":"user name","id":"user1","picture":"","profile":"","admin":false},"locator":{"site":"radio-t","url":"https://radio-t.com"},"score":0,"votes":{},"time":"2017-12-20T15:18:23-06:00"}` + "\n"
buf := &bytes.Buffer{}
buf.WriteString(r1)
buf.WriteString(r2)
buf.WriteString("{}")
inp := `{"version":1,"users":[{"id":"user1","blocked":{"status":false,"until":"0001-01-01T00:00:00Z"},"verified":true},{"id":"user2","blocked":{"status":true,"until":"2018-12-23T02:55:22.472041-06:00"},"verified":false}],"posts":[{"url":"https://radio-t.com","read_only":true}]}
{"id":"efbc17f177ee1a1c0ee6e1e025749966ec071adc","pid":"","text":"some text, <a href=\"http://radio-t.com\" rel=\"nofollow\">link</a>","user":{"name":"user name","id":"user1","picture":"","ip":"293ec5b0cf154855258824ec7fac5dc63d176915","admin":false},"locator":{"site":"radio-t","url":"https://radio-t.com"},"score":0,"votes":{},"time":"2017-12-20T15:18:22-06:00"}
{"id":"f863bd79-fec6-4a75-b308-61fe5dd02aa1","pid":"1234","text":"some text2","user":{"name":"user name","id":"user2","picture":"","ip":"293ec5b0cf154855258824ec7fac5dc63d176915","admin":false},"locator":{"site":"radio-t","url":"https://radio-t.com/2"},"score":0,"votes":{},"time":"2017-12-20T15:18:23-06:00"}`
b := prep(t) // write some recs
r := Native{DataStore: &service.DataStore{Interface: b, AdminStore: admin.NewStaticStore("12345", []string{}, "")}}
size, err := r.Import(buf, "radio-t")
size, err := r.Import(strings.NewReader(inp), "radio-t")
assert.Nil(t, err)
assert.Equal(t, 2, size)
comments, err := b.Find(store.Locator{SiteID: "radio-t", URL: "https://radio-t.com"}, "time")
comments, err := b.Last("radio-t", 10)
assert.Nil(t, err)
assert.Equal(t, 2, len(comments))
assert.Equal(t, "efbc17f177ee1a1c0ee6e1e025749966ec071adc", comments[0].ID)
assert.Equal(t, "afbc17f177ee1a1c0ee6e1e025749966ec071adc", comments[1].ID)
assert.Equal(t, "efbc17f177ee1a1c0ee6e1e025749966ec071adc", comments[1].ParentID)
assert.Equal(t, "f863bd79-fec6-4a75-b308-61fe5dd02aa1", comments[0].ID)
assert.Equal(t, "1234", comments[0].ParentID)
assert.Equal(t, false, b.IsReadOnly(comments[0].Locator))
// try import again
buf.WriteString(r1)
buf.WriteString(r2)
buf.WriteString("{}")
size, err = r.Import(buf, "radio-t")
assert.Nil(t, err)
assert.Equal(t, 2, size)
assert.Equal(t, "efbc17f177ee1a1c0ee6e1e025749966ec071adc", comments[1].ID)
assert.Equal(t, "https://radio-t.com", comments[1].Locator.URL)
assert.Equal(t, true, b.IsReadOnly(comments[1].Locator))
assert.Equal(t, false, b.IsBlocked("radio-t", "user1"))
assert.Equal(t, true, b.IsVerified("radio-t", "user1"))
assert.Equal(t, true, b.IsBlocked("radio-t", "user2"))
assert.Equal(t, false, b.IsVerified("radio-t", "user2"))
}
func TestNative_ImportManyWithError(t *testing.T) {
@@ -103,11 +106,12 @@ func TestNative_ImportManyWithError(t *testing.T) {
goodRec := `{"id":"%d","pid":"","text":"some text, <a href=\"http://radio-t.com\" rel=\"nofollow\">link</a>","user":{"name":"user name","id":"user1","picture":"","profile":"","admin":false},"locator":{"site":"radio-t","url":"https://radio-t.com"},"score":0,"votes":{},"time":"2017-12-20T15:18:22-06:00"}` + "\n"
buf := &bytes.Buffer{}
buf.WriteString(`{"version":1, "users":[], "posts":[]}` + "\n")
for i := 0; i < 1200; i++ {
buf.WriteString(fmt.Sprintf(goodRec, i))
}
buf.WriteString("bad1\n")
buf.WriteString("bad2\n")
buf.WriteString("{}\n")
buf.WriteString("{}\n")
b := prep(t) // write some recs
r := Native{DataStore: &service.DataStore{Interface: b, AdminStore: admin.NewStaticStore("12345", []string{}, "")}}
+26
View File
@@ -5,6 +5,8 @@ import (
"sync"
"time"
multierror "github.com/hashicorp/go-multierror"
"github.com/google/uuid"
"github.com/pkg/errors"
@@ -279,6 +281,30 @@ func (s *DataStore) Metas(siteID string) (umetas []UserMetaData, pmetas []PostMe
return umetas, pmetas, nil
}
// SetMetas saves metadata for users and posts
func (s *DataStore) SetMetas(siteID string, umetas []UserMetaData, pmetas []PostMetaData) (err error) {
errs := new(multierror.Error)
// save posts metas
for _, pm := range pmetas {
if pm.ReadOnly {
errs = multierror.Append(errs, s.SetReadOnly(store.Locator{SiteID: siteID, URL: pm.URL}, true))
}
}
// save users metas
for _, um := range umetas {
if um.Blocked.Status {
errs = multierror.Append(errs, s.SetBlock(siteID, um.ID, true, um.Blocked.Until.Sub(time.Now())))
}
if um.Verified {
errs = multierror.Append(errs, s.SetVerified(siteID, um.ID, true))
}
}
return errs.ErrorOrNil()
}
// getsScopedLocks pull lock from the map if found or create a new one
func (s *DataStore) getsScopedLocks(id string) (lock sync.Locker) {
s.scopedLocks.Do(func() { s.scopedLocks.locks = map[string]sync.Locker{} })