add image storage to memory_store

This commit is contained in:
Dmitry Verkhoturov
2020-04-04 14:06:39 -05:00
committed by Umputun
parent ce92f63215
commit 0355ba4ff7
9 changed files with 428 additions and 5 deletions
@@ -0,0 +1,105 @@
/*
* Copyright 2020 Umputun. All rights reserved.
* Use of this source code is governed by a MIT-style
* license that can be found in the LICENSE file.
*/
package accessor
import (
"context"
"path"
"sync"
"time"
log "github.com/go-pkgz/lgr"
"github.com/pkg/errors"
"github.com/rs/xid"
)
// MemImage implements image.Store with memory backend
type MemImage struct {
imagesStaging map[string][]byte
images map[string][]byte
insertTime map[string]time.Time
sync.RWMutex
}
// NewMemImageStore makes admin Store in memory.
func NewMemImageStore() *MemImage {
log.Print("[DEBUG] make memory image store")
return &MemImage{
imagesStaging: map[string][]byte{},
images: map[string][]byte{},
insertTime: map[string]time.Time{},
}
}
func (m *MemImage) Save(userID string, img []byte) (id string, err error) {
id = path.Join(userID, guid())
return m.SaveWithID(id, img)
}
func (m *MemImage) SaveWithID(id string, img []byte) (string, error) {
m.Lock()
m.imagesStaging[id] = img
m.insertTime[id] = time.Now()
m.Unlock()
return id, nil
}
func (m *MemImage) Load(id string) ([]byte, error) {
m.RLock()
img, ok := m.images[id]
if !ok {
img, ok = m.imagesStaging[id]
}
m.RUnlock()
if !ok {
return nil, errors.Errorf("image %s not found", id)
}
return img, nil
}
func (m *MemImage) Commit(id string) error {
m.RLock()
img, ok := m.imagesStaging[id]
m.RUnlock()
if !ok {
return errors.Errorf("failed to commit %s, not found in staging", id)
}
m.Lock()
m.images[id] = img
m.Unlock()
return nil
}
func (m *MemImage) Cleanup(_ context.Context, ttl time.Duration) error {
var idsToRemove []string
m.RLock()
for id, t := range m.insertTime {
age := time.Since(t)
if age > ttl {
log.Printf("[INFO] remove staging image %s, age %v", id, age)
idsToRemove = append(idsToRemove, id)
}
}
m.RUnlock()
m.Lock()
for _, id := range idsToRemove {
delete(m.insertTime, id)
delete(m.imagesStaging, id)
}
m.Unlock()
return nil
}
// guid makes a globally unique id
func guid() string {
return xid.New().String()
}
@@ -0,0 +1,95 @@
/*
* Copyright 2020 Umputun. All rights reserved.
* Use of this source code is governed by a MIT-style
* license that can be found in the LICENSE file.
*/
package accessor
import (
"context"
"encoding/base64"
"io"
"io/ioutil"
"strings"
"testing"
"time"
"github.com/stretchr/testify/assert"
)
// gopher png for test, from https://golang.org/src/image/png/example_test.go
const gopher = "iVBORw0KGgoAAAANSUhEUgAAAEsAAAA8CAAAAAALAhhPAAAFfUlEQVRYw62XeWwUVRzHf2" +
"+OPbo9d7tsWyiyaZti6eWGAhISoIGKECEKCAiJJkYTiUgTMYSIosYYBBIUIxoSPIINEBDi2VhwkQrVsj1ESgu9doHWdrul7ba" +
"73WNm3vOPtsseM9MdwvvrzTs+8/t95ze/33sI5BqiabU6m9En8oNjduLnAEDLUsQXFF8tQ5oxK3vmnNmDSMtrncks9Hhtt" +
"/qeWZapHb1ha3UqYSWVl2ZmpWgaXMXGohQAvmeop3bjTRtv6SgaK/Pb9/bFzUrYslbFAmHPp+3WhAYdr+7GN/YnpN46Opv55VDs" +
"JkoEpMrY/vO2BIYQ6LLvm0ThY3MzDzzeSJeeWNyTkgnIE5ePKsvKlcg/0T9QMzXalwXMlj54z4c0rh/mzEfr+FgWEz2w6uk" +
"8dkzFAgcARAgNp1ZYef8bH2AgvuStbc2/i6CiWGj98y2tw2l4FAXKkQBIf+exyRnteY83LfEwDQAYCoK+P6bxkZm/0966LxcAA" +
"ILHB56kgD95PPxltuYcMtFTWw/FKkY/6Opf3GGd9ZF+Qp6mzJxzuRSractOmJrH1u8XTvWFHINNkLQLMR+XHXvfPPHw967raE1xxwtA36I" +
"MRfkAAG29/7mLuQcb2WOnsJReZGfpiHsSBX81cvMKywYZHhX5hFPtOqPGWZCXnhWGAu6lX91ElKXSalcLXu3UaOXVay57ZSe5f6Gpx7J2" +
"MXAsi7EqSp09b/MirKSyJfnfEEgeDjl8FgDAfvewP03zZ+AJ0m9aFRM8eEHBDRKjfcreDXnZdQuAxXpT2NRJ7xl3UkLBhuVGU16gZiGOgZm" +
"rSbRdqkILuL/yYoSXHHkl9KXgqNu3PB8oRg0geC5vFmLjad6mUyTKLmF3OtraWDIfACyXqmephaDABawfpi6tqqBZytfQMqOz6S09iWXhkt" +
"rRaB8Xz4Yi/8gyABDm5NVe6qq/3VzPrcjELWrebVuyY2T7ar4zQyybUCtsQ5Es1FGaZVrRVQwAgHGW2ZCRZshI5bGQi7HesyE972pOSeMM0" +
"dSktlzxRdrlqb3Osa6CCS8IJoQQQgBAbTAa5l5epO34rJszibJI8rxLfGzcp1dRosutGeb2VDNgqYrwTiPNsLxXiPi3dz7LiS1WBRBDBOnqEj" +
"yy3aQb+/bLiJzz9dIkscVBBLxMfSEac7kO4Fpkngi0ruNBeSOal+u8jgOuqPz12nryMLCniEjtOOOmpt+KEIqsEdocJjYXwrh9OZqWJQyPCTo67" +
"LNS/TdxLAv6R5ZNK9npEjbYdT33gRo4o5oTqR34R+OmaSzDBWsAIPhuRcgyoteNi9gF0KzNYWVItPf2TLoXEg+7isNC7uJkgo1iQWOfRSP9NR" +
"11RtbZZ3OMG/VhL6jvx+J1m87+RCfJChAtEBQkSBX2PnSiihc/Twh3j0h7qdYQAoRVsRGmq7HU2QRbaxVGa1D6nIOqaIWRjyRZpHMQKWKpZM5fe" +
"A+lzC4ZFultV8S6T0mzQGhQohi5I8iw+CsqBSxhFMuwyLgSwbghGb0AiIKkSDmGZVmJSiKihsiyOAUs70UkywooYP0bii9GdH4sfr1UNysd3fU" +
"yLLMQN+rsmo3grHl9VNJHbbwxoa47Vw5gupIqrZcjPh9R4Nye3nRDk199V+aetmvVtDRE8/+cbgAAgMIWGb3UA0MGLE9SCbWX670TDy" +
"1y98c3D27eppUjsZ6fql3jcd5rUe7+ZIlLNQny3Rd+E5Tct3WVhTM5RBCEdiEK0b6B+/ca2gYU393nFj/n1AygRQxPIUA043M42u85+z2S" +
"nssKrPl8Mx76NL3E6eXc3be7OD+H4WHbJkKI8AU8irbITQjZ+0hQcPEgId/Fn/pl9crKH02+5o2b9T/eMx7pKoskYgAAAABJRU5ErkJggg=="
func gopherPNG() io.Reader { return base64.NewDecoder(base64.StdEncoding, strings.NewReader(gopher)) }
func TestMemImage_Save(t *testing.T) {
svc := NewMemImageStore()
id, err := svc.Save("user1", []byte(gopher))
assert.NoError(t, err)
assert.Contains(t, id, "user1/")
}
func TestMemImage_SaveWithIDFail(t *testing.T) {
svc := NewMemImageStore()
id, err := svc.SaveWithID("test_id", []byte(gopher))
assert.NoError(t, err)
assert.Equal(t, id, "test_id")
}
func TestMemImage_LoadAfterSave(t *testing.T) {
svc := NewMemImageStore()
gopher, err := ioutil.ReadAll(gopherPNG())
assert.NoError(t, err)
img, err := svc.Load("test_id")
assert.EqualError(t, err, "image test_id not found")
assert.Empty(t, img)
id, err := svc.Save("user1", gopher)
assert.NoError(t, err)
img, err = svc.Load(id)
assert.NoError(t, err)
assert.Equal(t, gopher, img)
err = svc.Commit(id)
assert.NoError(t, err)
err = svc.Cleanup(context.TODO(), 0)
assert.NoError(t, err)
img, err = svc.Load(id)
assert.NoError(t, err)
assert.Equal(t, gopher, img)
}
func TestMemImage_CommitFail(t *testing.T) {
svc := NewMemImageStore()
err := svc.Commit("test_id")
assert.EqualError(t, err, "failed to commit test_id, not found in staging")
}
func TestMemImage_Cleanup(t *testing.T) {
svc := NewMemImageStore()
err := svc.Cleanup(context.TODO(), time.Minute)
assert.NoError(t, err)
}
@@ -33,6 +33,8 @@ services:
- ADMIN_RPC_API=http://mem_store.r42:8080/cmd
- STORE_TYPE=rpc
- STORE_RPC_API=http://mem_store.r42:8080/cmd
- IMAGE_TYPE=rpc
- IMAGE_RPC_API=http://mem_store.r42:8080/cmd
mem_store.r42:
image: umputun/mem_store.r42
+1
View File
@@ -7,6 +7,7 @@ require (
github.com/go-pkgz/lgr v0.7.0
github.com/jessevdk/go-flags v1.4.0
github.com/pkg/errors v0.9.1
github.com/rs/xid v1.2.1
github.com/stretchr/testify v1.5.1
github.com/umputun/remark/backend v1.5.0
)
+3 -2
View File
@@ -1,5 +1,5 @@
/*
* Copyright 2019 Umputun. All rights reserved.
* Copyright 2020 Umputun. All rights reserved.
* Use of this source code is governed by a MIT-style
* license that can be found in the LICENSE file.
*/
@@ -41,6 +41,7 @@ func main() {
dataStore := accessor.NewMemData()
adminStore := accessor.NewMemAdminStore(opts.Secret)
imgStore := accessor.NewMemImageStore()
rpcServer := jrpc.Server{
API: opts.API,
@@ -51,7 +52,7 @@ func main() {
Logger: log.Default(),
}
srv := server.NewRPC(dataStore, adminStore, &rpcServer)
srv := server.NewRPC(dataStore, adminStore, imgStore, &rpcServer)
admRec := accessor.AdminRec{
SiteID: "remark",
@@ -0,0 +1,69 @@
/*
* Copyright 2020 Umputun. All rights reserved.
* Use of this source code is governed by a MIT-style
* license that can be found in the LICENSE file.
*/
package server
import (
"context"
"encoding/base64"
"encoding/json"
"time"
"github.com/go-pkgz/jrpc"
)
func (s *RPC) imgSaveHndl(id uint64, params json.RawMessage) (rr jrpc.Response) {
var req [2]string
if err := json.Unmarshal(params, &req); err != nil {
return jrpc.Response{Error: err.Error()}
}
img, err := base64.StdEncoding.DecodeString(req[1])
if err != nil {
return jrpc.Response{Error: err.Error()}
}
value, err := s.img.Save(req[0], img)
return jrpc.EncodeResponse(id, value, err)
}
func (s *RPC) imgSaveWithIDHndl(id uint64, params json.RawMessage) (rr jrpc.Response) {
var req [2]string
if err := json.Unmarshal(params, &req); err != nil {
return jrpc.Response{Error: err.Error()}
}
img, err := base64.StdEncoding.DecodeString(req[1])
if err != nil {
return jrpc.Response{Error: err.Error()}
}
value, err := s.img.SaveWithID(req[0], img)
return jrpc.EncodeResponse(id, value, err)
}
func (s *RPC) imgLoadHndl(id uint64, params json.RawMessage) (rr jrpc.Response) {
var fileID string
if err := json.Unmarshal(params, &fileID); err != nil {
return jrpc.Response{Error: err.Error()}
}
value, err := s.img.Load(fileID)
return jrpc.EncodeResponse(id, value, err)
}
func (s *RPC) imgCommitHndl(id uint64, params json.RawMessage) (rr jrpc.Response) {
var fileID string
if err := json.Unmarshal(params, &fileID); err != nil {
return jrpc.Response{Error: err.Error()}
}
err := s.img.Commit(fileID)
return jrpc.EncodeResponse(id, nil, err)
}
func (s *RPC) imgCleanupHndl(id uint64, params json.RawMessage) (rr jrpc.Response) {
var ttl time.Duration
if err := json.Unmarshal(params, &ttl); err != nil {
return jrpc.Response{Error: err.Error()}
}
err := s.img.Cleanup(context.TODO(), ttl)
return jrpc.EncodeResponse(id, nil, err)
}
@@ -0,0 +1,138 @@
/*
* Copyright 2020 Umputun. All rights reserved.
* Use of this source code is governed by a MIT-style
* license that can be found in the LICENSE file.
*/
package server
import (
"context"
"encoding/base64"
"fmt"
"io"
"io/ioutil"
"net/http"
"strings"
"testing"
"time"
"github.com/go-pkgz/jrpc"
"github.com/stretchr/testify/assert"
"github.com/umputun/remark/backend/app/store/image"
)
// gopher png for test, from https://golang.org/src/image/png/example_test.go
const gopher = "iVBORw0KGgoAAAANSUhEUgAAAEsAAAA8CAAAAAALAhhPAAAFfUlEQVRYw62XeWwUVRzHf2" +
"+OPbo9d7tsWyiyaZti6eWGAhISoIGKECEKCAiJJkYTiUgTMYSIosYYBBIUIxoSPIINEBDi2VhwkQrVsj1ESgu9doHWdrul7ba" +
"73WNm3vOPtsseM9MdwvvrzTs+8/t95ze/33sI5BqiabU6m9En8oNjduLnAEDLUsQXFF8tQ5oxK3vmnNmDSMtrncks9Hhtt" +
"/qeWZapHb1ha3UqYSWVl2ZmpWgaXMXGohQAvmeop3bjTRtv6SgaK/Pb9/bFzUrYslbFAmHPp+3WhAYdr+7GN/YnpN46Opv55VDs" +
"JkoEpMrY/vO2BIYQ6LLvm0ThY3MzDzzeSJeeWNyTkgnIE5ePKsvKlcg/0T9QMzXalwXMlj54z4c0rh/mzEfr+FgWEz2w6uk" +
"8dkzFAgcARAgNp1ZYef8bH2AgvuStbc2/i6CiWGj98y2tw2l4FAXKkQBIf+exyRnteY83LfEwDQAYCoK+P6bxkZm/0966LxcAA" +
"ILHB56kgD95PPxltuYcMtFTWw/FKkY/6Opf3GGd9ZF+Qp6mzJxzuRSractOmJrH1u8XTvWFHINNkLQLMR+XHXvfPPHw967raE1xxwtA36I" +
"MRfkAAG29/7mLuQcb2WOnsJReZGfpiHsSBX81cvMKywYZHhX5hFPtOqPGWZCXnhWGAu6lX91ElKXSalcLXu3UaOXVay57ZSe5f6Gpx7J2" +
"MXAsi7EqSp09b/MirKSyJfnfEEgeDjl8FgDAfvewP03zZ+AJ0m9aFRM8eEHBDRKjfcreDXnZdQuAxXpT2NRJ7xl3UkLBhuVGU16gZiGOgZm" +
"rSbRdqkILuL/yYoSXHHkl9KXgqNu3PB8oRg0geC5vFmLjad6mUyTKLmF3OtraWDIfACyXqmephaDABawfpi6tqqBZytfQMqOz6S09iWXhkt" +
"rRaB8Xz4Yi/8gyABDm5NVe6qq/3VzPrcjELWrebVuyY2T7ar4zQyybUCtsQ5Es1FGaZVrRVQwAgHGW2ZCRZshI5bGQi7HesyE972pOSeMM0" +
"dSktlzxRdrlqb3Osa6CCS8IJoQQQgBAbTAa5l5epO34rJszibJI8rxLfGzcp1dRosutGeb2VDNgqYrwTiPNsLxXiPi3dz7LiS1WBRBDBOnqEj" +
"yy3aQb+/bLiJzz9dIkscVBBLxMfSEac7kO4Fpkngi0ruNBeSOal+u8jgOuqPz12nryMLCniEjtOOOmpt+KEIqsEdocJjYXwrh9OZqWJQyPCTo67" +
"LNS/TdxLAv6R5ZNK9npEjbYdT33gRo4o5oTqR34R+OmaSzDBWsAIPhuRcgyoteNi9gF0KzNYWVItPf2TLoXEg+7isNC7uJkgo1iQWOfRSP9NR" +
"11RtbZZ3OMG/VhL6jvx+J1m87+RCfJChAtEBQkSBX2PnSiihc/Twh3j0h7qdYQAoRVsRGmq7HU2QRbaxVGa1D6nIOqaIWRjyRZpHMQKWKpZM5fe" +
"A+lzC4ZFultV8S6T0mzQGhQohi5I8iw+CsqBSxhFMuwyLgSwbghGb0AiIKkSDmGZVmJSiKihsiyOAUs70UkywooYP0bii9GdH4sfr1UNysd3fU" +
"yLLMQN+rsmo3grHl9VNJHbbwxoa47Vw5gupIqrZcjPh9R4Nye3nRDk199V+aetmvVtDRE8/+cbgAAgMIWGb3UA0MGLE9SCbWX670TDy" +
"1y98c3D27eppUjsZ6fql3jcd5rUe7+ZIlLNQny3Rd+E5Tct3WVhTM5RBCEdiEK0b6B+/ca2gYU393nFj/n1AygRQxPIUA043M42u85+z2S" +
"nssKrPl8Mx76NL3E6eXc3be7OD+H4WHbJkKI8AU8irbITQjZ+0hQcPEgId/Fn/pl9crKH02+5o2b9T/eMx7pKoskYgAAAABJRU5ErkJggg=="
func gopherPNG() io.Reader { return base64.NewDecoder(base64.StdEncoding, strings.NewReader(gopher)) }
func gopherPNGBytes() []byte {
img, _ := ioutil.ReadAll(gopherPNG())
return img
}
func TestRPC_imgSaveHndl(t *testing.T) {
_, port, teardown := prepTestStore(t)
defer teardown()
api := fmt.Sprintf("http://localhost:%d/test", port)
ri := image.RPC{Client: jrpc.Client{API: api, Client: http.Client{Timeout: 1 * time.Second}}}
id, err := ri.Save("admin", gopherPNGBytes())
assert.NoError(t, err)
assert.Contains(t, id, "admin/", "id contains username")
err = ri.Commit(id)
assert.NoError(t, err)
}
func TestRPC_imgSaveWithIDHndl(t *testing.T) {
_, port, teardown := prepTestStore(t)
defer teardown()
api := fmt.Sprintf("http://localhost:%d/test", port)
ri := image.RPC{Client: jrpc.Client{API: api, Client: http.Client{Timeout: 1 * time.Second}}}
id, err := ri.SaveWithID("test_id", gopherPNGBytes())
assert.NoError(t, err)
assert.Equal(t, id, "test_id")
}
func TestRPC_imgLoadHndl(t *testing.T) {
_, port, teardown := prepTestStore(t)
defer teardown()
api := fmt.Sprintf("http://localhost:%d/test", port)
ri := image.RPC{Client: jrpc.Client{API: api, Client: http.Client{Timeout: 1 * time.Second}}}
// save
id, err := ri.Save("admin", gopherPNGBytes())
assert.NoError(t, err)
// load
img, err := ri.Load(id)
assert.NoError(t, err)
assert.Equal(t, 1462, len(img))
assert.Equal(t, gopherPNGBytes(), img)
// commit
err = ri.Commit(id)
assert.NoError(t, err)
// load after commit
img, err = ri.Load(id)
assert.NoError(t, err)
assert.Equal(t, 1462, len(img))
assert.Equal(t, gopherPNGBytes(), img)
}
func TestRPC_imgCommitHndlFail(t *testing.T) {
_, port, teardown := prepTestStore(t)
defer teardown()
api := fmt.Sprintf("http://localhost:%d/test", port)
ri := image.RPC{Client: jrpc.Client{API: api, Client: http.Client{Timeout: 1 * time.Second}}}
err := ri.Commit("test_id")
assert.EqualError(t, err, "failed to commit test_id, not found in staging")
}
func TestRPC_imgCleanupHndl(t *testing.T) {
_, port, teardown := prepTestStore(t)
defer teardown()
api := fmt.Sprintf("http://localhost:%d/test", port)
ri := image.RPC{Client: jrpc.Client{API: api, Client: http.Client{Timeout: 1 * time.Second}}}
// save
id, err := ri.Save("admin", gopherPNGBytes())
assert.NoError(t, err)
// load
_, err = ri.Load(id)
assert.NoError(t, err)
// cleanup
err = ri.Cleanup(context.TODO(), time.Nanosecond)
assert.NoError(t, err)
// load after cleanup should fail
_, err = ri.Load(id)
assert.Error(t, err)
assert.Contains(t, err.Error(), "image admin/")
assert.Contains(t, err.Error(), "not found")
}
+13 -2
View File
@@ -8,6 +8,7 @@ package server
import (
"github.com/go-pkgz/jrpc"
"github.com/umputun/remark/backend/app/store/image"
"github.com/umputun/remark/backend/app/store/admin"
"github.com/umputun/remark/backend/app/store/engine"
@@ -19,11 +20,12 @@ type RPC struct {
*jrpc.Server
eng engine.Interface
adm admin.Store
img image.Store
}
// NewRPC makes RPC instance and register handlers
func NewRPC(e engine.Interface, a admin.Store, r *jrpc.Server) *RPC {
res := &RPC{eng: e, adm: a, Server: r}
func NewRPC(e engine.Interface, a admin.Store, i image.Store, r *jrpc.Server) *RPC {
res := &RPC{eng: e, adm: a, img: i, Server: r}
res.addHandlers()
return res
}
@@ -52,4 +54,13 @@ func (s *RPC) addHandlers() {
"enabled": s.admEnabledHndl,
"event": s.admEventHndl,
})
// image store handlers
s.Group("image", jrpc.HandlersGroup{
"save": s.imgSaveHndl,
"save_with_id": s.imgSaveWithIDHndl,
"load": s.imgLoadHndl,
"commit": s.imgCommitHndl,
"cleanup": s.imgCleanupHndl,
})
}
@@ -47,7 +47,8 @@ func waitForHTTPServerStart(port int) {
func prepTestStore(t *testing.T) (s *RPC, port int, teardown func()) {
mg := accessor.NewMemData()
adm := accessor.NewMemAdminStore("secret")
s = NewRPC(mg, adm, &jrpc.Server{API: "/test", Logger: jrpc.NoOpLogger})
img := accessor.NewMemImageStore()
s = NewRPC(mg, adm, img, &jrpc.Server{API: "/test", Logger: jrpc.NoOpLogger})
admRec := accessor.AdminRec{
SiteID: "test-site",