add image.Store.Info() function
This commit is contained in:
@@ -13,6 +13,8 @@ import (
|
||||
|
||||
log "github.com/go-pkgz/lgr"
|
||||
"github.com/pkg/errors"
|
||||
|
||||
"github.com/umputun/remark/backend/app/store/image"
|
||||
)
|
||||
|
||||
// MemImage implements image.Store with memory backend
|
||||
@@ -95,3 +97,17 @@ func (m *MemImage) Cleanup(_ context.Context, ttl time.Duration) error {
|
||||
m.Unlock()
|
||||
return nil
|
||||
}
|
||||
|
||||
// Info returns meta information about storage
|
||||
func (m *MemImage) Info() (image.StoreInfo, error) {
|
||||
var ts time.Time
|
||||
m.RLock()
|
||||
for _, t := range m.insertTime {
|
||||
if ts.IsZero() || t.Before(ts) {
|
||||
ts = t
|
||||
}
|
||||
}
|
||||
m.RUnlock()
|
||||
|
||||
return image.StoreInfo{FirstStagingImageTS: ts}, nil
|
||||
}
|
||||
|
||||
@@ -80,3 +80,23 @@ func TestMemImage_Cleanup(t *testing.T) {
|
||||
err := svc.Cleanup(context.TODO(), time.Minute)
|
||||
assert.NoError(t, err)
|
||||
}
|
||||
|
||||
func TestMemImage_Info(t *testing.T) {
|
||||
svc := NewMemImageStore()
|
||||
gopher, err := ioutil.ReadAll(gopherPNG())
|
||||
assert.NoError(t, err)
|
||||
|
||||
// get info on empty storage, should be zero
|
||||
info, err := svc.Info()
|
||||
assert.NoError(t, err)
|
||||
assert.True(t, info.FirstStagingImageTS.IsZero())
|
||||
|
||||
// save image
|
||||
err = svc.Save("test_img", gopher)
|
||||
assert.NoError(t, err)
|
||||
|
||||
// get info after saving, should be non-zero
|
||||
info, err = svc.Info()
|
||||
assert.NoError(t, err)
|
||||
assert.False(t, info.FirstStagingImageTS.IsZero())
|
||||
}
|
||||
|
||||
@@ -54,3 +54,8 @@ func (s *RPC) imgCleanupHndl(id uint64, params json.RawMessage) (rr jrpc.Respons
|
||||
err := s.img.Cleanup(context.TODO(), ttl)
|
||||
return jrpc.EncodeResponse(id, nil, err)
|
||||
}
|
||||
|
||||
func (s *RPC) imgInfoHndl(id uint64, _ json.RawMessage) (rr jrpc.Response) {
|
||||
info, err := s.img.Info()
|
||||
return jrpc.EncodeResponse(id, info, err)
|
||||
}
|
||||
|
||||
@@ -124,3 +124,25 @@ func TestRPC_imgCleanupHndl(t *testing.T) {
|
||||
_, err = ri.Load(id)
|
||||
assert.EqualError(t, err, "image test_img not found")
|
||||
}
|
||||
|
||||
func TestRPC_imgInfoHndl(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}}}
|
||||
|
||||
// get info on empty storage, should be zero
|
||||
info, err := ri.Info()
|
||||
assert.NoError(t, err)
|
||||
assert.True(t, info.FirstStagingImageTS.IsZero())
|
||||
|
||||
// save
|
||||
err = ri.Save("test_img", gopherPNGBytes())
|
||||
assert.NoError(t, err)
|
||||
|
||||
// get info after saving, should be non-zero
|
||||
info, err = ri.Info()
|
||||
assert.NoError(t, err)
|
||||
assert.False(t, info.FirstStagingImageTS.IsZero())
|
||||
}
|
||||
|
||||
@@ -61,5 +61,6 @@ func (s *RPC) addHandlers() {
|
||||
"load": s.imgLoadHndl,
|
||||
"commit": s.imgCommitHndl,
|
||||
"cleanup": s.imgCleanupHndl,
|
||||
"info": s.imgInfoHndl,
|
||||
})
|
||||
}
|
||||
|
||||
@@ -140,3 +140,26 @@ func (b *Bolt) Cleanup(_ context.Context, ttl time.Duration) error {
|
||||
})
|
||||
return err
|
||||
}
|
||||
|
||||
// Info returns meta information about storage
|
||||
func (b *Bolt) Info() (StoreInfo, error) {
|
||||
var ts time.Time
|
||||
err := b.db.View(func(tx *bolt.Tx) error {
|
||||
c := tx.Bucket([]byte(insertTimeBktName)).Cursor()
|
||||
|
||||
for id, tsData := c.First(); id != nil; id, tsData = c.Next() {
|
||||
var createdRaw int64
|
||||
err := binary.Read(bytes.NewReader(tsData), binary.LittleEndian, &createdRaw)
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "failed to deserialize timestamp for %s", id)
|
||||
}
|
||||
|
||||
created := time.Unix(0, createdRaw)
|
||||
if ts.IsZero() || created.Before(ts) {
|
||||
ts = created
|
||||
}
|
||||
}
|
||||
return nil
|
||||
})
|
||||
return StoreInfo{FirstStagingImageTS: ts}, errors.Wrapf(err, "problem retrieving first timestamp from staging images")
|
||||
}
|
||||
|
||||
@@ -105,6 +105,25 @@ func TestBoltStore_Cleanup(t *testing.T) {
|
||||
assert.NoError(t, err)
|
||||
}
|
||||
|
||||
func TestBolt_Info(t *testing.T) {
|
||||
svc, teardown := prepareBoltImageStorageTest(t)
|
||||
defer teardown()
|
||||
|
||||
// get info on empty storage, should be zero
|
||||
info, err := svc.Info()
|
||||
assert.NoError(t, err)
|
||||
assert.True(t, info.FirstStagingImageTS.IsZero())
|
||||
|
||||
// save image
|
||||
err = svc.Save("test_img", gopherPNGBytes())
|
||||
assert.NoError(t, err)
|
||||
|
||||
// get info after saving, should be non-zero
|
||||
info, err = svc.Info()
|
||||
assert.NoError(t, err)
|
||||
assert.False(t, info.FirstStagingImageTS.IsZero())
|
||||
}
|
||||
|
||||
func assertBoltImgNil(t *testing.T, db *bolt.DB, bucket string, id string) {
|
||||
checkBoltImgData(t, db, bucket, id, func(data []byte) error {
|
||||
assert.Nil(t, data, id)
|
||||
|
||||
@@ -115,6 +115,33 @@ func (f *FileSystem) Cleanup(_ context.Context, ttl time.Duration) error {
|
||||
return errors.Wrap(err, "failed to cleanup images")
|
||||
}
|
||||
|
||||
// Info returns meta information about storage
|
||||
func (f *FileSystem) Info() (StoreInfo, error) {
|
||||
if _, err := os.Stat(f.Staging); os.IsNotExist(err) {
|
||||
return StoreInfo{}, nil
|
||||
}
|
||||
|
||||
var ts time.Time
|
||||
err := filepath.Walk(f.Staging, func(fpath string, info os.FileInfo, err error) error {
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if info.IsDir() {
|
||||
return nil
|
||||
}
|
||||
|
||||
created := info.ModTime()
|
||||
if ts.IsZero() || created.Before(ts) {
|
||||
ts = created
|
||||
}
|
||||
return nil
|
||||
})
|
||||
if err != nil {
|
||||
return StoreInfo{}, errors.Wrapf(err, "problem retrieving first timestamp from staging images on fs")
|
||||
}
|
||||
return StoreInfo{FirstStagingImageTS: ts}, nil
|
||||
}
|
||||
|
||||
// location gets full path for id by adding partition to the final path in order to keep files in different subdirectories
|
||||
// and avoid too many files in a single place.
|
||||
// the end result is a full path like this - /tmp/images/user1/92/xxx-yyy.png.
|
||||
|
||||
@@ -107,7 +107,6 @@ func TestFsStore_LoadAfterSave(t *testing.T) {
|
||||
id := "test_img"
|
||||
err := svc.Save(id, gopherPNGBytes())
|
||||
assert.NoError(t, err)
|
||||
t.Log(id)
|
||||
|
||||
data, err := svc.Load(id)
|
||||
assert.NoError(t, err)
|
||||
@@ -125,7 +124,6 @@ func TestFsStore_LoadAfterCommit(t *testing.T) {
|
||||
id := "test_img"
|
||||
err := svc.Save(id, gopherPNGBytes())
|
||||
assert.NoError(t, err)
|
||||
t.Log(id)
|
||||
err = svc.Commit(id)
|
||||
require.NoError(t, err)
|
||||
|
||||
@@ -233,6 +231,25 @@ func TestFsStore_Cleanup(t *testing.T) {
|
||||
assert.Error(t, err, "no file on staging anymore")
|
||||
}
|
||||
|
||||
func TestFsStore_Info(t *testing.T) {
|
||||
svc, teardown := prepareImageTest(t)
|
||||
defer teardown()
|
||||
|
||||
// get ts on empty storage, should be zero
|
||||
ts, err := svc.Info()
|
||||
assert.NoError(t, err)
|
||||
assert.True(t, ts.FirstStagingImageTS.IsZero())
|
||||
|
||||
// save image
|
||||
err = svc.Save("test_img", gopherPNGBytes())
|
||||
assert.NoError(t, err)
|
||||
|
||||
// get ts after saving, should be non-zero
|
||||
ts, err = svc.Info()
|
||||
assert.NoError(t, err)
|
||||
assert.False(t, ts.FirstStagingImageTS.IsZero())
|
||||
}
|
||||
|
||||
func prepareImageTest(t *testing.T) (svc *FileSystem, teardown func()) {
|
||||
loc, err := ioutil.TempDir("", "test_image_r42")
|
||||
require.NoError(t, err, "failed to make temp dir")
|
||||
|
||||
@@ -52,6 +52,11 @@ type ServiceParams struct {
|
||||
MaxWidth int
|
||||
}
|
||||
|
||||
// StoreInfo contains image store meta information
|
||||
type StoreInfo struct {
|
||||
FirstStagingImageTS time.Time
|
||||
}
|
||||
|
||||
// To regenerate mock run from this directory:
|
||||
// sh -c "mockery -inpkg -name Store -print > /tmp/image-mock.tmp && mv /tmp/image-mock.tmp image_mock.go"
|
||||
|
||||
@@ -60,6 +65,7 @@ type ServiceParams struct {
|
||||
// Two-stage commit scheme is used for not storing images which are uploaded but later never used in the comments,
|
||||
// e.g. when somebody uploaded a picture but did not sent the comment.
|
||||
type Store interface {
|
||||
Info() (StoreInfo, error) // get meta information about storage
|
||||
Save(id string, img []byte) error // store image with passed id to staging
|
||||
Load(id string) ([]byte, error) // load image by ID
|
||||
|
||||
@@ -150,6 +156,11 @@ func (s *Service) Cleanup(ctx context.Context) {
|
||||
}
|
||||
}
|
||||
|
||||
// Info returns meta information about storage
|
||||
func (s *Service) Info() (StoreInfo, error) {
|
||||
return s.store.Info()
|
||||
}
|
||||
|
||||
// Close flushes all in-progress submits and enforces waiting commits
|
||||
func (s *Service) Close(ctx context.Context) {
|
||||
log.Printf("[INFO] close image service ")
|
||||
|
||||
@@ -39,6 +39,27 @@ func (_m *MockStore) Commit(id string) error {
|
||||
return r0
|
||||
}
|
||||
|
||||
// Info provides a mock function with given fields:
|
||||
func (_m *MockStore) Info() (StoreInfo, error) {
|
||||
ret := _m.Called()
|
||||
|
||||
var r0 StoreInfo
|
||||
if rf, ok := ret.Get(0).(func() StoreInfo); ok {
|
||||
r0 = rf()
|
||||
} else {
|
||||
r0 = ret.Get(0).(StoreInfo)
|
||||
}
|
||||
|
||||
var r1 error
|
||||
if rf, ok := ret.Get(1).(func() error); ok {
|
||||
r1 = rf()
|
||||
} else {
|
||||
r1 = ret.Error(1)
|
||||
}
|
||||
|
||||
return r0, r1
|
||||
}
|
||||
|
||||
// Load provides a mock function with given fields: id
|
||||
func (_m *MockStore) Load(id string) ([]byte, error) {
|
||||
ret := _m.Called(id)
|
||||
|
||||
@@ -140,6 +140,17 @@ func TestService_SubmitDelay(t *testing.T) {
|
||||
store.AssertNumberOfCalls(t, "Commit", 5)
|
||||
}
|
||||
|
||||
func TestService_Info(t *testing.T) {
|
||||
store := MockStore{}
|
||||
store.On("Info", mock.Anything, mock.Anything).Once().Return(StoreInfo{}, nil)
|
||||
|
||||
svc := Service{store: &store, ServiceParams: ServiceParams{}}
|
||||
info, err := svc.Info()
|
||||
assert.NoError(t, err)
|
||||
assert.True(t, info.FirstStagingImageTS.IsZero())
|
||||
store.AssertNumberOfCalls(t, "Info", 1)
|
||||
}
|
||||
|
||||
func TestService_resize(t *testing.T) {
|
||||
// reader is nil
|
||||
resized := resize(nil, 100, 100)
|
||||
|
||||
@@ -46,3 +46,16 @@ func (r *RPC) Cleanup(_ context.Context, ttl time.Duration) error {
|
||||
_, err := r.Call("image.cleanup", ttl)
|
||||
return err
|
||||
}
|
||||
|
||||
// Info returns meta information about storage
|
||||
func (r *RPC) Info() (StoreInfo, error) {
|
||||
resp, err := r.Call("image.info")
|
||||
if err != nil {
|
||||
return StoreInfo{}, err
|
||||
}
|
||||
info := StoreInfo{}
|
||||
if err = json.Unmarshal(*resp.Result, &info); err != nil {
|
||||
return StoreInfo{}, err
|
||||
}
|
||||
return info, err
|
||||
}
|
||||
|
||||
@@ -65,6 +65,20 @@ func TestRemote_Cleanup(t *testing.T) {
|
||||
assert.NoError(t, err)
|
||||
}
|
||||
|
||||
func TestRemote_Info(t *testing.T) {
|
||||
ts := testServer(t, `{"method":"image.info","id":1}`,
|
||||
`{"result":{"FirstStagingImageTS":"0001-01-01T00:00:01Z"},"id":1}`)
|
||||
defer ts.Close()
|
||||
c := RPC{Client: jrpc.Client{API: ts.URL, Client: http.Client{}}}
|
||||
|
||||
var a Store = &c
|
||||
_ = a
|
||||
|
||||
info, err := c.Info()
|
||||
assert.NoError(t, err)
|
||||
assert.False(t, info.FirstStagingImageTS.IsZero())
|
||||
}
|
||||
|
||||
func testServer(t *testing.T, req, resp string) *httptest.Server {
|
||||
return httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
body, err := ioutil.ReadAll(r.Body)
|
||||
|
||||
Reference in New Issue
Block a user