Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
37 commits
Select commit Hold shift + click to select a range
4d02031
add vendor dir, simple build.sh, route package
Apr 18, 2018
c22c823
add space to align
Apr 18, 2018
e2ec1ff
Merge branch 'develop' into feature/routing
Apr 18, 2018
8fadd6b
remove .gitkeep
Apr 18, 2018
736656c
add consistent hash
Apr 18, 2018
c4628fc
enlarge ring namespace to uint64 instead of uint32, use fnv64-a inste…
Apr 20, 2018
f5ab651
fix file name typo :-p
Apr 20, 2018
5c172b5
hash ring can store NodeInfo, all UT passed, Cheers
Apr 20, 2018
a7c2a40
DHT by consistent hash RPC works now, need to implement more API
Apr 25, 2018
c11ea34
format code, need some tools to do this automaticly
Apr 25, 2018
e63071c
add .travis.yml for CI
Apr 25, 2018
f955958
fix build
Apr 25, 2018
d6012c2
add reviewdog at Travis
Apr 25, 2018
41d765e
add dep install
Apr 25, 2018
3d9d6b3
fix missing goverage golint for Travis
Apr 25, 2018
0204c57
fix build.sh use default GOOS GOARCH
Apr 25, 2018
d76a6a0
update Travis key
Apr 25, 2018
ea85542
for reviewdog test
Apr 25, 2018
b03602a
run test and goverage
Apr 25, 2018
757fa30
fix .travis.yml
Apr 25, 2018
ee047cf
fix .travis.yml
Apr 25, 2018
b014d17
fix travis.yml .....
Apr 25, 2018
3644774
fix travis
Apr 25, 2018
3216c2f
try reviewdog
Apr 26, 2018
edf646f
make golint happy
Apr 26, 2018
f4daa08
fix golint issue, add necessary comment
Apr 26, 2018
ce2adce
fix UT issue, typo
Apr 26, 2018
db9c714
use thunderdb api token
Apr 26, 2018
5c238e6
add graphic dep analysis
Apr 27, 2018
67fe296
add analysisVendor.sh to generate dep tree png
May 2, 2018
a18a760
add test case
May 2, 2018
0179cb7
add Envelope to RPC proto
May 2, 2018
f36541c
add new vendor
May 2, 2018
6bac220
add ignore
May 2, 2018
2b53cd1
Merge branch 'develop' of github.com:thunderdb/ThunderDB into develop
May 2, 2018
e643ce2
Merge branch 'develop' into feature/routing
May 2, 2018
313f331
ignore 'vendor/' 'server/' 'utils/' on golint
May 2, 2018
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
@@ -1,6 +1,8 @@
coverage.txt
analysisVendor.png

# Binaries for programs and plugins
bin/
*.exe
*.dll
*.so
Expand Down
4 changes: 2 additions & 2 deletions .travis.yml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
env:
global:
- secure: "dld6R6RF3Mtjs94cBj94eg1phi/34objqKlZI64A+mI68J5yXTxoLRJJMZrzM63YQ8wbvvF+W68KDyjXvYaQXTFES/p+igx5UsQ3SQ7GtRuqjXkCWI0hVqRphFWeeWPWEactFmwIK/hap2dvcBIP0sIU100EE/Cx6Lo0FFkmMb664wA+by4ELFo0l3h7pPfklyVAAWY+GYp9jeLT8/7VT7aDqASUh+3znd7NBgokXkCAwOVNTKI5H0eHIyW5aM9vBKJvTuTTpZ51rAA7tUUMOvihh2Hltmxy6rBfyApUCfMDWTrdKnhttXKDuHvL2vh+GzTD1cl6HAo+owYP6YkGPuA40FfmgBB8DtoAIvCY0jU0hwkkZAh2uXHjthNUFDOGY26ZkuL+l/T5TxQEqe3WgIG2MoVUCrc9L+XqsmTuM3ERS3gdusmjSc/WwQDpy3DLkc5W2zokcgmolmGnaFkxqwifsCTK3Ifzr8yKsSWZHXQadRRI43dmxua3b4WNFAY/lsF6yrbRyWY9LtWwTjbdheoMOIH6mE1daHIXKN1T5xPwPK4EW99Y1jDdJHdpWc2zYHnkEm8H9CJ9+KuAMADB9813XDBkcXRtvLtJckhdGPmkKrqJ3WHnMGT/D4ErMWPAY6+oawruizReDmvi3KA3Emp2VDfleAe0pCsyDwwqCKw="
- secure: "CNZXODvXHMB0z6Z+Jd8MjzP+RYgRzkykRNpZdCDCo6AC57BpV5U20Mn9o2MGhTA9lMgtyNXXEaQxeIGZMFp3iB2UHiPif1UD4C0uDttytTuCcmHMlBIVP8GiNcAexi1Zv2SchM1k8z3tw5zDhkYMwjs/4ZsGVXHa1ip/UzpZaJBPl9D64wedJxZGBjS/Kn9IZTrMkL2otr8rKk0ZFMtwDkHKdjfikZ+PAlkYpX2a2Dp69HTLbEtx+ugocNj02hzgVYNpFXA1K0CulyCYJq7oqdx0UFbtKEd3urYFWfmxEZSg6SoHcIujC+/BKfZELFyA3uomp+cLpGaRlEp/iAlj/ixTo4iFtu6XdDOim4J9ZVtNSg6DdDstL9uhxCu+vca/ntViGfyw6CF0hKcHTZvlH2awiV1njmHFdRvoaMy/1fdw1JAV9O7ylU3bV5GrF6fMQ1CpweYCnF3IMpvvwdeUnGR6NPvqjybcnyXoTXZ139PxfWzbUzSG8KbOvzt1aVhWF5184b8fqhN3768Y1stke+FjFpARqlPVsfzS7fnn5kx8zZ5nfmwZjnL/WFob24mvm+/9e7IsfWmeEQdFspD+spBE4WOta7tegqh5ANNfM2owd9evtH5kqlCLeAVD3LMs69ZQW45KS0rDJpY8+ffMWRJRmoM/pweW3CmbHPeb468="
- REVIEWDOG_VERSION=0.9.8
language: go
go:
Expand Down Expand Up @@ -32,4 +32,4 @@ script:
- go test -v -race $(go list ./... | grep -v "/vendor/")
- goverage -coverprofile=coverage.txt ./...
- >-
golint ./... | reviewdog -f=golint -ci=travis
golint ./... | grep -v 'vendor/' | grep -v 'server/' | grep -v 'utils/' | reviewdog -f=golint -ci=travis
4 changes: 4 additions & 0 deletions analysisVendor.sh
Original file line number Diff line number Diff line change
@@ -0,0 +1,4 @@
#!/bin/sh

dep ensure
dep status -dot | dot -Tpng -o analysisVendor.png && open analysisVendor.png
88 changes: 57 additions & 31 deletions consistent/consistent_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,9 +5,7 @@
package consistent

import (
"bufio"
"math/rand"
"os"
"runtime"
"sort"
"strconv"
Expand Down Expand Up @@ -56,6 +54,15 @@ func TestRemove(t *testing.T) {
utils.CheckNum(len(x.sortedHashes), 0, t)
}

func TestMembers(t *testing.T) {
x := New()
x.Add(NewNodeFromID("abcdefg"))
x.Add(NewNodeFromID("abcdefghi"))
x.Remove(NewNodeFromID("abcdefg"))
utils.CheckNum(len(x.Members()), 1, t)
utils.CheckStr(string(x.Members()[0].ID), "abcdefghi", t)
}

func TestRemoveNonExisting(t *testing.T) {
x := New()
x.Add(NewNodeFromID("abcdefg"))
Expand Down Expand Up @@ -215,6 +222,14 @@ func TestGetTwo(t *testing.T) {
}
}

func TestGetTwoEmpty(t *testing.T) {
x := New()
_, _, err := x.GetTwo("9999999")
if err != ErrEmptyCircle {
t.Fatal(err)
}
}

func TestGetTwoQuick(t *testing.T) {
x := New()
x.Add(NewNodeFromID("abcdefg"))
Expand Down Expand Up @@ -360,6 +375,17 @@ func TestGetNMore(t *testing.T) {
}
}

func TestGetNEmpty(t *testing.T) {
x := New()
members, err := x.GetN("9999999", 5)
if err != ErrEmptyCircle {
t.Fatal(err)
}
if len(members) != 0 {
t.Errorf("expected 3 members instead of %d", len(members))
}
}

func TestGetNQuick(t *testing.T) {
x := New()
x.Add(NewNodeFromID("abcdefg"))
Expand Down Expand Up @@ -675,35 +701,35 @@ func TestAddCollision(t *testing.T) {
}
}

// inspired by @or-else on github
func TestCollisionsCRC(t *testing.T) {
t.SkipNow()
c := New()
f, err := os.Open("/usr/share/dict/words")
if err != nil {
t.Fatal(err)
}
defer f.Close()
found := make(map[NodeKey]string)
scanner := bufio.NewScanner(f)
count := 0
for scanner.Scan() {
word := scanner.Text()
for i := 0; i < c.NumberOfReplicas; i++ {
ekey := c.nodeKey(NodeID(word), i)
// ekey := word + "|" + strconv.Itoa(i)
k := c.hashKey(ekey)
exist, ok := found[k]
if ok {
t.Logf("found collision: %v, %v", ekey, exist)
count++
} else {
found[k] = ekey
}
}
}
t.Logf("number of collisions: %d", count)
}
//// inspired by @or-else on github
//func TestCollisionsCRC(t *testing.T) {
// t.SkipNow()
// c := New()
// f, err := os.Open("/usr/share/dict/words")
// if err != nil {
// t.Fatal(err)
// }
// defer f.Close()
// found := make(map[NodeKey]string)
// scanner := bufio.NewScanner(f)
// count := 0
// for scanner.Scan() {
// word := scanner.Text()
// for i := 0; i < c.NumberOfReplicas; i++ {
// ekey := c.nodeKey(NodeID(word), i)
// // ekey := word + "|" + strconv.Itoa(i)
// k := c.hashKey(ekey)
// exist, ok := found[k]
// if ok {
// t.Logf("found collision: %v, %v", ekey, exist)
// count++
// } else {
// found[k] = ekey
// }
// }
// }
// t.Logf("number of collisions: %d", count)
//}

func TestConcurrentGetSet(t *testing.T) {
x := New()
Expand Down
41 changes: 41 additions & 0 deletions proto/proto.go
Original file line number Diff line number Diff line change
@@ -1,5 +1,9 @@
package proto

import (
"time"
)

// NodeID is node name, will be generated from Hash(nodePublicKey)
type NodeID string

Expand All @@ -14,3 +18,40 @@ type Node struct {
ID NodeID
PublicKey string
}

// Envelope is the protocol
type Envelope struct {
Version string
TTL time.Duration
Expire time.Duration
}

// PingReq is Ping RPC request
type PingReq struct {
Node Node
Version string
Envelope
}

// PingResp is Ping RPC response, i.e. Pong
type PingResp struct {
Msg string
Version string
Envelope
}

// FindValueReq is FindValue RPC request
type FindValueReq struct {
NodeID NodeID
Count int
Version string
Envelope
}

// FindValueResp is FindValue RPC response
type FindValueResp struct {
Nodes []Node
Msg string
Version string
Envelope
}
43 changes: 31 additions & 12 deletions route/dht_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,14 +2,14 @@ package route

import (
log "github.com/sirupsen/logrus"
"github.com/thunderdb/ThunderDB/proto"
. "github.com/smartystreets/goconvey/convey"
. "github.com/thunderdb/ThunderDB/proto"
"github.com/thunderdb/ThunderDB/rpc"
"github.com/thunderdb/ThunderDB/utils"
"net"
"testing"
)

func TestGetNeighbors(t *testing.T) {
func TestPingFindValue(t *testing.T) {
log.SetLevel(log.DebugLevel)
addr := "127.0.0.1:0"
l, err := net.Listen("tcp", addr)
Expand All @@ -29,29 +29,48 @@ func TestGetNeighbors(t *testing.T) {
log.Fatal(err)
}

reqA := &AddNodeReq{
Node: proto.Node{
reqA := &PingReq{
Node: Node{
ID: "node1",
},
}
respA := new(AddNodeResp)
err = client.Call("DHT.AddNode", reqA, respA)
respA := new(PingResp)
err = client.Call("DHT.Ping", reqA, respA)
if err != nil {
log.Fatal(err)
}
log.Debugf("respA: %v", respA)

req := &GetNeighborsReq{
nodeID: "123",
reqB := &PingReq{
Node: Node{
ID: "node2",
},
}
respB := new(PingResp)
err = client.Call("DHT.Ping", reqB, respB)
if err != nil {
log.Fatal(err)
}
log.Debugf("respA: %v", respB)

req := &FindValueReq{
NodeID: "123",
Count: 2,
}
resp := new(GetNeighborsResp)
err = client.Call("DHT.GetNeighbors", req, resp)
resp := new(FindValueResp)
err = client.Call("DHT.FindValue", req, resp)
if err != nil {
log.Fatal(err)
}
log.Debugf("resp: %v", resp)
utils.CheckStr(string(resp.Nodes[0].ID), "node1", t)
nodeIDList := []string{
string(resp.Nodes[0].ID),
string(resp.Nodes[1].ID),
}
Convey("test FindValue", t, func() {
So(nodeIDList, ShouldContain, "node1")
So(nodeIDList, ShouldContain, "node2")
})

client.Close()
dhtServer.Stop()
Expand Down
42 changes: 9 additions & 33 deletions route/service.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@ package route

import (
"fmt"
"github.com/opentracing/opentracing-go/log"
log "github.com/sirupsen/logrus"
"github.com/thunderdb/ThunderDB/consistent"
"github.com/thunderdb/ThunderDB/proto"
)
Expand All @@ -20,44 +20,20 @@ func NewDHTService() *DHTService {
}
}

// GetNeighborsReq is GetNeighbors RPC request
type GetNeighborsReq struct {
nodeID proto.NodeID
Count int
Version string
}

// GetNeighborsResp is GetNeighbors RPC response
type GetNeighborsResp struct {
Nodes []proto.Node
ErrMsg string
Version string
}

// GetNeighbors RPC returns GetNeighborsReq.Count closest node from consistent hash ring
func (DHT *DHTService) GetNeighbors(req *GetNeighborsReq, resp *GetNeighborsResp) (err error) {
resp.Nodes, err = DHT.hashRing.GetN(string(req.nodeID), req.Count)
// FindValue RPC returns FindValueReq.Count closest node from DHT
func (DHT *DHTService) FindValue(req *proto.FindValueReq, resp *proto.FindValueResp) (err error) {
resp.Nodes, err = DHT.hashRing.GetN(string(req.NodeID), req.Count)
if err != nil {
log.Error(err)
resp.ErrMsg = fmt.Sprint(err)
resp.Msg = fmt.Sprint(err)
}
return
}

// AddNodeReq is AddNode RPC request
type AddNodeReq struct {
Node proto.Node
Version string
}

// AddNodeResp is AddNode RPC response
type AddNodeResp struct {
ErrMsg string
Version string
}

// AddNode RPC add AddNodeReq.Node to consistent hash ring
func (DHT *DHTService) AddNode(req *AddNodeReq, resp *AddNodeResp) (err error) {
// Ping RPC add PingReq.Node to DHT
func (DHT *DHTService) Ping(req *proto.PingReq, resp *proto.PingResp) (err error) {
DHT.hashRing.Add(req.Node)
resp = new(proto.PingResp)
resp.Msg = "Pong"
return
}
24 changes: 24 additions & 0 deletions vendor/github.com/gopherjs/gopherjs/LICENSE

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Loading