@cryptotaxi247 / kubo / commits / 999910626

feat: initial redis datastore implementation

Brian Tiger Chow committed Feb 6, 2015 at 15:48 UTC 99991062678020df391d3dec93a6a2b6c507e582
2 files changed +192
thirdparty/redis-datastore/datastore.go new
+84
@@ -0,0 +1,84 @@
1 +package redis
2 +
3 +import (
4 + "errors"
5 + "fmt"
6 + "sync"
7 + "time"
8 +
9 + "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/fzzy/radix/redis"
10 + datastore "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore"
11 + query "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore/query"
12 +)
13 +
14 +var _ datastore.Datastore = &RedisDatastore{}
15 +var _ datastore.ThreadSafeDatastore = &RedisDatastore{}
16 +
17 +var ErrInvalidType = errors.New("redis datastore: invalid type error. this datastore only supports []byte values")
18 +
19 +func NewExpiringDatastore(client *redis.Client, ttl time.Duration) (datastore.ThreadSafeDatastore, error) {
20 + return &RedisDatastore{
21 + client: client,
22 + ttl: ttl,
23 + }, nil
24 +}
25 +
26 +func NewDatastore(client *redis.Client) (datastore.ThreadSafeDatastore, error) {
27 + return &RedisDatastore{
28 + client: client,
29 + }, nil
30 +}
31 +
32 +type RedisDatastore struct {
33 + mu sync.Mutex
34 + client *redis.Client
35 + ttl time.Duration
36 +}
37 +
38 +func (ds *RedisDatastore) Put(key datastore.Key, value interface{}) error {
39 + ds.mu.Lock()
40 + defer ds.mu.Unlock()
41 +
42 + data, ok := value.([]byte)
43 + if !ok {
44 + return ErrInvalidType
45 + }
46 +
47 + ds.client.Append("SET", key.String(), data)
48 + if ds.ttl != 0 {
49 + ds.client.Append("EXPIRE", key.String(), ds.ttl.Seconds())
50 + }
51 + if err := ds.client.GetReply().Err; err != nil {
52 + return fmt.Errorf("failed to put value: %s", err)
53 + }
54 + if ds.ttl != 0 {
55 + if err := ds.client.GetReply().Err; err != nil {
56 + return fmt.Errorf("failed to set expiration: %s", err)
57 + }
58 + }
59 + return nil
60 +}
61 +
62 +func (ds *RedisDatastore) Get(key datastore.Key) (value interface{}, err error) {
63 + ds.mu.Lock()
64 + defer ds.mu.Unlock()
65 + return ds.client.Cmd("GET", key.String()).Bytes()
66 +}
67 +
68 +func (ds *RedisDatastore) Has(key datastore.Key) (exists bool, err error) {
69 + ds.mu.Lock()
70 + defer ds.mu.Unlock()
71 + return ds.client.Cmd("EXISTS", key.String()).Bool()
72 +}
73 +
74 +func (ds *RedisDatastore) Delete(key datastore.Key) (err error) {
75 + ds.mu.Lock()
76 + defer ds.mu.Unlock()
77 + return ds.client.Cmd("DEL", key.String()).Err
78 +}
79 +
80 +func (ds *RedisDatastore) Query(q query.Query) (query.Results, error) {
81 + return nil, errors.New("TODO implement query for redis datastore?")
82 +}
83 +
84 +func (ds *RedisDatastore) IsThreadSafe() {}
thirdparty/redis-datastore/datastore_test.go new
+108
@@ -0,0 +1,108 @@
1 +package redis
2 +
3 +import (
4 + "bytes"
5 + "os"
6 + "testing"
7 + "time"
8 +
9 + "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/fzzy/radix/redis"
10 + datastore "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore"
11 + "github.com/jbenet/go-ipfs/thirdparty/assert"
12 +)
13 +
14 +const RedisEnv = "REDIS_DATASTORE_TEST_HOST"
15 +
16 +func TestPutGetBytes(t *testing.T) {
17 + client := clientOrAbort(t)
18 + ds, err := NewDatastore(client)
19 + if err != nil {
20 + t.Fatal(err)
21 + }
22 + key, val := datastore.NewKey("foo"), []byte("bar")
23 + assert.Nil(ds.Put(key, val), t)
24 + v, err := ds.Get(key)
25 + if err != nil {
26 + t.Fatal(err)
27 + }
28 + if bytes.Compare(v.([]byte), val) != 0 {
29 + t.Fail()
30 + }
31 +}
32 +
33 +func TestHasBytes(t *testing.T) {
34 + client := clientOrAbort(t)
35 + ds, err := NewDatastore(client)
36 + if err != nil {
37 + t.Fatal(err)
38 + }
39 + key, val := datastore.NewKey("foo"), []byte("bar")
40 + has, err := ds.Has(key)
41 + if err != nil {
42 + t.Fatal(err)
43 + }
44 + if has {
45 + t.Fail()
46 + }
47 +
48 + assert.Nil(ds.Put(key, val), t)
49 + hasAfterPut, err := ds.Has(key)
50 + if err != nil {
51 + t.Fatal(err)
52 + }
53 + if !hasAfterPut {
54 + t.Fail()
55 + }
56 +}
57 +
58 +func TestDelete(t *testing.T) {
59 + client := clientOrAbort(t)
60 + ds, err := NewDatastore(client)
61 + if err != nil {
62 + t.Fatal(err)
63 + }
64 + key, val := datastore.NewKey("foo"), []byte("bar")
65 + assert.Nil(ds.Put(key, val), t)
66 + assert.Nil(ds.Delete(key), t)
67 +
68 + hasAfterDelete, err := ds.Has(key)
69 + if err != nil {
70 + t.Fatal(err)
71 + }
72 + if hasAfterDelete {
73 + t.Fail()
74 + }
75 +}
76 +
77 +func TestExpiry(t *testing.T) {
78 + ttl := 1 * time.Second
79 + client := clientOrAbort(t)
80 + ds, err := NewExpiringDatastore(client, ttl)
81 + if err != nil {
82 + t.Fatal(err)
83 + }
84 + key, val := datastore.NewKey("foo"), []byte("bar")
85 + assert.Nil(ds.Put(key, val), t)
86 + time.Sleep(ttl + 1*time.Second)
87 + assert.Nil(ds.Delete(key), t)
88 +
89 + hasAfterExpiration, err := ds.Has(key)
90 + if err != nil {
91 + t.Fatal(err)
92 + }
93 + if hasAfterExpiration {
94 + t.Fail()
95 + }
96 +}
97 +
98 +func clientOrAbort(t *testing.T) *redis.Client {
99 + c, err := redis.Dial("tcp", os.Getenv(RedisEnv))
100 + if err != nil {
101 + t.Log("could not connect to a redis instance")
102 + t.SkipNow()
103 + }
104 + if err := c.Cmd("FLUSHALL").Err; err != nil {
105 + t.Fatal(err)
106 + }
107 + return c
108 +}