diff --git a/core/batch.go b/core/batch.go index eb15732c4..5030f3f18 100644 --- a/core/batch.go +++ b/core/batch.go @@ -3,25 +3,23 @@ package core import ( "encoding/gob" "pilosa/db" - . "pilosa/util" - - "github.com/gocql/gocql/uuid" + "pilosa/util" ) type BatchRequest struct { - Id *uuid.UUID - Source *uuid.UUID - Fragment_id SUUID + Id *util.GUID + Source *util.GUID + Fragment_id util.SUUID Bitmap_id uint64 Compressed_bitmap string Filter uint64 } type BatchResponse struct { - Id *uuid.UUID + Id *util.GUID } -func (self BatchResponse) ResultId() *uuid.UUID { +func (self BatchResponse) ResultId() *util.GUID { return self.Id } func (self BatchResponse) ResultData() interface{} { @@ -41,7 +39,7 @@ func (self *Service) Batch(database_name, frame, compressed_bitmap string, bitma fragment, err := database.GetFragmentForBitmap(oslice, &db.Bitmap{bitmap_id, frame, filter}) if err == nil { - id := uuid.RandomUUID() + id := util.RandomUUID() batch := db.Message{Data: BatchRequest{Id: &id, Source: self.Id, Fragment_id: fragment.GetId(), Bitmap_id: bitmap_id, Compressed_bitmap: compressed_bitmap}} dest_id := fragment.GetProcess().Id() self.Transport.Send(&batch, &dest_id) diff --git a/core/etcd.go b/core/etcd.go index ca5ca19c0..a9f143132 100644 --- a/core/etcd.go +++ b/core/etcd.go @@ -14,7 +14,6 @@ import ( "github.com/coreos/go-etcd/etcd" "github.com/davecgh/go-spew/spew" - "github.com/gocql/gocql/uuid" ) type TopologyMapper struct { @@ -46,9 +45,14 @@ func (self *TopologyMapper) Run() { receiver := make(chan *etcd.Response) stop := make(chan bool) go func() { - // TODO: error check and restart watcher + // TODO: add some terminating measure // TODO: use modindex to make sure watch catches everything - _, _ = self.service.Etcd.Watch(self.namespace+"/db", 0, true, receiver, stop) + for { + ns := self.namespace + "/db" + log.Println(" ETCD watcher:", ns) + resp, err = self.service.Etcd.Watch(ns, 0, true, receiver, stop) + log.Println("TopologyMapper ETCD watcher", resp, err) + } }() go func() { for resp = range receiver { @@ -77,15 +81,32 @@ func (p PairList) Swap(i, j int) { p[i], p[j] = p[j], p[i] } func (p PairList) Len() int { return len(p) } func (p PairList) Less(i, j int) bool { return p[i].Value < p[j].Value } -func getLightestProcess(m map[string]int) Pair { +func getLightestProcess(m map[string]int) (Pair, error) { - p := make(PairList, len(m)) + l := len(m) + if l == 0 { + return Pair{}, errors.New("No Processes") + } + p := make(PairList, l) i := 0 for k, v := range m { p[i] = Pair{k, v} } sort.Sort(p) - return p[0] + + return p[l-1], nil +} + +func (self *TopologyMapper) MakeFragments(db string, slice_int int) error { + frames_to_create := config.GetStringArrayDefault("supported_frames", []string{"b.n", "l.n", "t.t", "d"}) + for _, frame := range frames_to_create { + err := self.AllocateFragment(db, frame, slice_int) + if err != nil { + log.Println(err) + } + } + return nil + } func (self *TopologyMapper) AllocateFragment(db, frame string, slice_int int) error { @@ -114,16 +135,19 @@ func (self *TopologyMapper) AllocateFragment(db, frame string, slice_int int) er } } - p := getLightestProcess(m) + p, err := getLightestProcess(m) + if err != nil { + return err + } // PUT -d "value=5cb315c3-6e1d-4218-89b7-943d1dba985b" http://etcd0:4001/v2/keys/pilosa/0/db/29/frame/d/slice/5/fragment/a2b632fc4001b817/proces //so i need db, frame, slice , fragment_id fuid := util.SUUID_to_Hex(util.Id()) - fragment_key := fmt.Sprintf("%s/db/%d/frame/%s/slice/%d/fragment/%s/process", self.namespace, db, frame, slice_int, fuid) + fragment_key := fmt.Sprintf("%s/db/%s/frame/%s/slice/%d/fragment/%s/process", self.namespace, db, frame, slice_int, fuid) process_guid := p.Key // need to check value to see how many we have left _, err = self.service.Etcd.Set(fragment_key, process_guid, 0) - log.Println("Fragment sent to etcd: %s(%s)", fragment_key, process_guid) + log.Printf("Fragment sent to etcd: %s(%s)", fragment_key, process_guid) case 400: //key already present return errors.New("Fragment creation already in process:" + lock_key) default: @@ -145,7 +169,7 @@ func (self *TopologyMapper) handlenode(node *etcd.Node) error { var fragment_id util.SUUID var slice *db.Slice var slice_int int - var process_uuid uuid.UUID + var process_uuid util.GUID var process *db.Process var err error @@ -189,7 +213,7 @@ func (self *TopologyMapper) handlenode(node *etcd.Node) error { if bits[8] != "process" { return errors.New("no process") } - process_uuid, err = uuid.ParseUUID(node.Value) + process_uuid, err = util.ParseGUID(node.Value) if err != nil { return err } @@ -207,26 +231,26 @@ func (self *TopologyMapper) handlenode(node *etcd.Node) error { func flatten(node *etcd.Node) []*etcd.Node { nodes := []*etcd.Node{node} for i := 0; i < len(node.Nodes); i++ { - nodes = append(nodes, flatten(&node.Nodes[i])...) + nodes = append(nodes, flatten(node.Nodes[i])...) } return nodes } type Node struct { - id *uuid.UUID + id *util.GUID ip string port_tcp int port_http int } type ProcessMap struct { - nodes map[uuid.UUID]*db.Process + nodes map[util.GUID]*db.Process mutex sync.Mutex } func NewProcessMap() *ProcessMap { p := ProcessMap{} - p.nodes = make(map[uuid.UUID]*db.Process) + p.nodes = make(map[util.GUID]*db.Process) return &p } @@ -236,7 +260,7 @@ func (self *ProcessMap) AddProcess(process *db.Process) { self.nodes[process.Id()] = process } -func (self *ProcessMap) GetProcess(id *uuid.UUID) (*db.Process, error) { +func (self *ProcessMap) GetProcess(id *util.GUID) (*db.Process, error) { self.mutex.Lock() defer self.mutex.Unlock() process, ok := self.nodes[*id] @@ -246,7 +270,7 @@ func (self *ProcessMap) GetProcess(id *uuid.UUID) (*db.Process, error) { return process, nil } -func (self *ProcessMap) GetOrAddProcess(id *uuid.UUID) *db.Process { +func (self *ProcessMap) GetOrAddProcess(id *util.GUID) *db.Process { process, err := self.GetProcess(id) if err != nil { process = db.NewProcess(id) @@ -255,7 +279,7 @@ func (self *ProcessMap) GetOrAddProcess(id *uuid.UUID) *db.Process { return process } -func (self *ProcessMap) GetHost(id *uuid.UUID) (string, error) { +func (self *ProcessMap) GetHost(id *util.GUID) (string, error) { self.mutex.Lock() defer self.mutex.Unlock() process, ok := self.nodes[*id] @@ -265,7 +289,7 @@ func (self *ProcessMap) GetHost(id *uuid.UUID) (string, error) { return process.Host(), nil } -func (self *ProcessMap) GetPortTcp(id *uuid.UUID) (int, error) { +func (self *ProcessMap) GetPortTcp(id *util.GUID) (int, error) { self.mutex.Lock() defer self.mutex.Unlock() process, ok := self.nodes[*id] @@ -275,7 +299,7 @@ func (self *ProcessMap) GetPortTcp(id *uuid.UUID) (int, error) { return process.PortTcp(), nil } -func (self *ProcessMap) GetPortHttp(id *uuid.UUID) (int, error) { +func (self *ProcessMap) GetPortHttp(id *util.GUID) (int, error) { self.mutex.Lock() defer self.mutex.Unlock() process, ok := self.nodes[*id] @@ -324,7 +348,7 @@ func getKey(input string) string { return bits[len(bits)-1] } -func (self *ProcessMapper) getnode(u *uuid.UUID) *Node { +func (self *ProcessMapper) getnode(u *util.GUID) *Node { return new(Node) } @@ -344,9 +368,9 @@ func (self *ProcessMapper) handlenode(node *etcd.Node) error { } if len(bits) >= 2 { id_string := bits[1] - id, err := uuid.ParseUUID(id_string) + id, err := util.ParseGUID(id_string) if err != nil { - return errors.New("Invalid UUID: " + id_string) + return errors.New("Invalid GUID: " + id_string) } process = self.service.ProcessMap.GetOrAddProcess(&id) } diff --git a/core/http.go b/core/http.go index 0fd97101f..b2dce0316 100644 --- a/core/http.go +++ b/core/http.go @@ -22,7 +22,6 @@ import ( notify "github.com/bitly/go-notify" "github.com/davecgh/go-spew/spew" - "github.com/gocql/gocql/uuid" "github.com/gorilla/websocket" ) @@ -484,7 +483,7 @@ func (self *WebService) HandlePing(w http.ResponseWriter, r *http.Request) { return } process_string := r.Form.Get("process") - process_id, err := uuid.ParseUUID(process_string) + process_id, err := util.ParseGUID(process_string) if err != nil { http.Error(w, err.Error(), http.StatusBadRequest) return diff --git a/core/ping.go b/core/ping.go index 634111aa2..62fd3efb5 100644 --- a/core/ping.go +++ b/core/ping.go @@ -3,21 +3,20 @@ package core import ( "encoding/gob" "pilosa/db" + "pilosa/util" "time" - - "github.com/gocql/gocql/uuid" ) type PingRequest struct { - Id *uuid.UUID - Source *uuid.UUID + Id *util.GUID + Source *util.GUID } type PongRequest struct { - Id *uuid.UUID + Id *util.GUID } -func (self PongRequest) ResultId() *uuid.UUID { +func (self PongRequest) ResultId() *util.GUID { return self.Id } func (self PongRequest) ResultData() interface{} { @@ -29,8 +28,8 @@ func init() { gob.Register(PongRequest{}) } -func (self *Service) Ping(process_id *uuid.UUID) (*time.Duration, error) { - id := uuid.RandomUUID() +func (self *Service) Ping(process_id *util.GUID) (*time.Duration, error) { + id := util.RandomUUID() ping := db.Message{Data: PingRequest{Id: &id, Source: self.Id}} start := time.Now() self.Transport.Send(&ping, process_id) diff --git a/core/service.go b/core/service.go index f1a7dc6e8..317908426 100644 --- a/core/service.go +++ b/core/service.go @@ -14,12 +14,11 @@ import ( "syscall" "github.com/coreos/go-etcd/etcd" - "github.com/gocql/gocql/uuid" ) type Service struct { Stopper - Id *uuid.UUID + Id *util.GUID Etcd *etcd.Client Cluster *db.Cluster TopologyMapper *TopologyMapper @@ -71,17 +70,17 @@ func (self *Service) PrepareLogging() { } func (service *Service) init_id() { - var id uuid.UUID + var id util.GUID var err error id_string := config.GetString("id") if id_string == "" { log.Println("Service id not configured, generating...") - id = uuid.RandomUUID() + id = util.RandomUUID() if err != nil { log.Fatal("problem generating uuid") } } else { - id, err = uuid.ParseUUID(id_string) + id, err = util.ParseGUID(id_string) if err != nil { log.Fatalf("Service id '%s' not valid", id_string) } diff --git a/db/db.go b/db/db.go index d0708c8be..11489b0e9 100644 --- a/db/db.go +++ b/db/db.go @@ -2,8 +2,7 @@ package db import ( "encoding/gob" - - "github.com/gocql/gocql/uuid" + . "pilosa/util" ) type Message struct { @@ -12,11 +11,11 @@ type Message struct { type Envelope struct { Message *Message - Host *uuid.UUID + Host *GUID } type HoldResult interface { - ResultId() *uuid.UUID + ResultId() *GUID ResultData() interface{} } diff --git a/db/topology.go b/db/topology.go index 0f0cd909e..591247d72 100644 --- a/db/topology.go +++ b/db/topology.go @@ -7,7 +7,6 @@ import ( "pilosa/util" "sync" - "github.com/gocql/gocql/uuid" "github.com/stathat/consistent" ) @@ -17,23 +16,23 @@ var FragmentDoesNotExistError = errors.New("Fragment does not exist.") var FrameSliceIntersectDoesNotExistError = errors.New("FrameSliceIntersect does not exist.") type Location struct { - ProcessId *uuid.UUID + ProcessId *util.GUID FragmentId util.SUUID } type Process struct { - id *uuid.UUID + id *util.GUID host string port_tcp int port_http int mutex sync.Mutex } -func NewProcess(id *uuid.UUID) *Process { +func NewProcess(id *util.GUID) *Process { return &Process{id: id} } -func (self *Process) Id() uuid.UUID { +func (self *Process) Id() util.GUID { self.mutex.Lock() defer self.mutex.Unlock() return *self.id @@ -214,6 +213,10 @@ type Slice struct { id int } +func (self *Slice) Id() int { + return self.id +} + // Get a slice from a database func (d *Database) getSlice(slice_id int) (*Slice, error) { for _, slice := range d.slices { @@ -306,7 +309,7 @@ func (self *Fragment) GetProcess() *Process { return self.process } -func (self *Fragment) GetProcessId() *uuid.UUID { +func (self *Fragment) GetProcessId() *util.GUID { return self.process.id } @@ -342,7 +345,7 @@ func (d *Database) GetFragmentForBitmap(slice *Slice, bitmap *Bitmap) (*Fragment /* // NOT IMPLEMENTED // this would loop through all frame_slice_intersect[], then all fragmments to find a match -func (d *Database) GetFragmentById(fragment_id *uuid.UUID) *Fragment { +func (d *Database) GetFragmentById(fragment_id *GUID) *Fragment { } */ diff --git a/executor/executor.go b/executor/executor.go index 7711f66d4..966895aff 100644 --- a/executor/executor.go +++ b/executor/executor.go @@ -73,6 +73,10 @@ func (self *Executor) runQuery(database *db.Database, qry *query.Query) error { query_plan, err := query.QueryPlanForQuery(database, qry, &destination) if err != nil { + obj, found := err.(*query.FragmentNotFound) + if found { + self.service.TopologyMapper.MakeFragments(obj.Db, obj.Slice) + } self.service.Hold.Set(qry.Id, err, 30) return err } diff --git a/hold/hold.go b/hold/hold.go index 295d95bea..660d44978 100644 --- a/hold/hold.go +++ b/hold/hold.go @@ -2,41 +2,40 @@ package hold import ( "errors" + . "pilosa/util" "time" - - "github.com/gocql/gocql/uuid" ) type holdchan chan interface{} type gethold struct { - id *uuid.UUID + id *GUID reply chan holdchan } type delhold struct { - id *uuid.UUID + id *GUID } type Holder struct { - data map[uuid.UUID]holdchan + data map[GUID]holdchan getchan chan gethold delchan chan delhold } //var Hold Holder -func (self *Holder) DelChan(id *uuid.UUID) { +func (self *Holder) DelChan(id *GUID) { req := delhold{id} self.delchan <- req } -func (self *Holder) GetChan(id *uuid.UUID) holdchan { +func (self *Holder) GetChan(id *GUID) holdchan { reply := make(chan holdchan) req := gethold{id, reply} self.getchan <- req return <-reply } -func (self *Holder) Get(id *uuid.UUID, timeout int) (interface{}, error) { +func (self *Holder) Get(id *GUID, timeout int) (interface{}, error) { ch := self.GetChan(id) select { case val := <-ch: @@ -47,7 +46,7 @@ func (self *Holder) Get(id *uuid.UUID, timeout int) (interface{}, error) { } } -func (self *Holder) Set(id *uuid.UUID, value interface{}, timeout int) { +func (self *Holder) Set(id *GUID, value interface{}, timeout int) { ch := self.GetChan(id) go func() { select { @@ -77,13 +76,13 @@ func (self *Holder) Run() { } func NewHolder() *Holder { - h := Holder{make(map[uuid.UUID]holdchan), make(chan gethold), make(chan delhold)} + h := Holder{make(map[GUID]holdchan), make(chan gethold), make(chan delhold)} return &h } /* func init() { - Hold = Holder{make(map[uuid.UUID]holdchan), make(chan gethold), make(chan delhold)} + Hold = Holder{make(map[GUID]holdchan), make(chan gethold), make(chan delhold)} go Hold.Run() } */ diff --git a/hold/hold_test.go b/hold/hold_test.go index fa3f03d21..a16894df2 100644 --- a/hold/hold_test.go +++ b/hold/hold_test.go @@ -11,17 +11,17 @@ import ( func TestHoldChan(t *testing.T) { - Hold := Holder{make(map[uuid.UUID]holdchan), make(chan gethold), make(chan delhold)} + Hold := Holder{make(map[GUID]holdchan), make(chan gethold), make(chan delhold)} go Hold.Run() Convey("set then get", t, func() { - id := uuid.RandomUUID() + id := util.RandomUUID() Hold.Set(&id, "derp", 10) derp, _ := Hold.Get(&id, 10) So(derp, ShouldEqual, "derp") }) Convey("get then set", t, func() { - id := uuid.RandomUUID() + id := util.RandomUUID() go func() { Hold.Set(&id, "derpsy", 10) }() diff --git a/index/fragment_container.go b/index/fragment_container.go index 7253e00c6..af36a0306 100644 --- a/index/fragment_container.go +++ b/index/fragment_container.go @@ -202,7 +202,7 @@ func (self *FragmentContainer) Clear(frag_id SUUID) (bool, error) { } func (self *FragmentContainer) AddFragment(db string, frame string, slice int, id SUUID) { - log.Println("ADD FRAGMENT", frame) + log.Println("ADD FRAGMENT", frame, db, slice) f := NewFragment(id, db, slice, frame) self.fragments[id] = f diff --git a/index/general.go b/index/general.go index 95b899bf6..4842b29fa 100644 --- a/index/general.go +++ b/index/general.go @@ -101,8 +101,15 @@ func (self *General) Persist() error { defer w.Close() defer self.storage.Close() + results := make([]uint64, len(self.keys)) + i := 0 + for k, _ := range self.keys { // map[uint64]*Rank + results[i] = k + i += 1 + } + encoder := json.NewEncoder(w) - return encoder.Encode(self.keys) + return encoder.Encode(results) } func (self *General) Load(requestChan chan Command, f *Fragment) { @@ -114,12 +121,13 @@ func (self *General) Load(requestChan chan Command, f *Fragment) { } dec := json.NewDecoder(r) - var keys map[uint64]interface{} + var keys []uint64 + if err := dec.Decode(&keys); err != nil { return //log.Println("Bad mojo") } - for k, _ := range keys { + for _, k := range keys { request := NewLoadRequest(k) requestChan <- request request.Response() diff --git a/interfaces/core.go b/interfaces/core.go index cd6677f05..87e2e6b03 100644 --- a/interfaces/core.go +++ b/interfaces/core.go @@ -2,14 +2,13 @@ package interfaces import ( "pilosa/db" - - "github.com/gocql/gocql/uuid" + "pilosa/util" ) type Transporter interface { Run() Close() - Send(*db.Message, *uuid.UUID) + Send(*db.Message, *util.GUID) Receive() *db.Message Push(*db.Message) } diff --git a/query/parser.go b/query/parser.go index d5d1e2818..3c9f1a8d7 100644 --- a/query/parser.go +++ b/query/parser.go @@ -3,10 +3,10 @@ package query import ( "errors" "fmt" + "pilosa/util" "strconv" "github.com/davecgh/go-spew/spew" - "github.com/gocql/gocql/uuid" ) var InvalidQueryError = errors.New("Invalid query format.") @@ -66,7 +66,7 @@ func (self *QueryParser) Parse() (query *Query, err error) { }() var token *Token - id := uuid.RandomUUID() + id := util.RandomUUID() query = &Query{Id: &id, Subqueries: make([]Query, 0), Args: make(map[string]interface{})} token = self.next() diff --git a/query/planner.go b/query/planner.go index 6883e58e5..cda044851 100644 --- a/query/planner.go +++ b/query/planner.go @@ -2,29 +2,44 @@ package query import ( "encoding/gob" + "fmt" "math/rand" "pilosa/db" - - "github.com/gocql/gocql/uuid" + "pilosa/util" ) type PortableQueryStep interface { - GetId() *uuid.UUID + GetId() *util.GUID GetLocation() *db.Location } +type FragmentNotFound struct { + Db string + Frame string + Slice int + Retry bool +} + +func NewFragmentNotFound(db, frame string, slice int) *FragmentNotFound { + return &FragmentNotFound{db, frame, slice, true} +} + +func (self *FragmentNotFound) Error() string { + return fmt.Sprintf("Fragment Not Found: %s:%s:%d", self.Db, self.Frame, self.Slice) +} + /////////////////////////////////////////////////////////////////////////////////////////////////// // BASE /////////////////////////////////////////////////////////////////////////////////////////////////// type BaseQueryStep struct { - Id *uuid.UUID + Id *util.GUID Operation string Location *db.Location Destination *db.Location } -func (self *BaseQueryStep) GetId() *uuid.UUID { +func (self *BaseQueryStep) GetId() *util.GUID { return self.Id } func (self *BaseQueryStep) GetLocation() *db.Location { @@ -39,11 +54,11 @@ func (self *BaseQueryStep) LocIsDest() bool { } type BaseQueryResult struct { - Id *uuid.UUID + Id *util.GUID Data interface{} } -func (self *BaseQueryResult) ResultId() *uuid.UUID { +func (self *BaseQueryResult) ResultId() *util.GUID { return self.Id } @@ -56,7 +71,7 @@ func (self *BaseQueryResult) ResultData() interface{} { /////////////////////////////////////////////////////////////////////////////////////////////////// type CountQueryStep struct { *BaseQueryStep - Input *uuid.UUID + Input *util.GUID } type CountQueryResult struct { @@ -79,7 +94,7 @@ func (qt *CountQueryTree) getLocation(d *db.Database) (*db.Location, error) { /////////////////////////////////////////////////////////////////////////////////////////////////// type TopNQueryStep struct { *BaseQueryStep - Input *uuid.UUID + Input *util.GUID Filters []uint64 N int } @@ -106,7 +121,7 @@ func (qt *TopNQueryTree) getLocation(d *db.Database) (*db.Location, error) { /////////////////////////////////////////////////////////////////////////////////////////////////// type UnionQueryStep struct { *BaseQueryStep - Inputs []*uuid.UUID + Inputs []*util.GUID } type UnionQueryResult struct { @@ -138,7 +153,7 @@ func (qt *UnionQueryTree) getLocation(d *db.Database) (*db.Location, error) { /////////////////////////////////////////////////////////////////////////////////////////////////// type IntersectQueryStep struct { *BaseQueryStep - Inputs []*uuid.UUID + Inputs []*util.GUID } type IntersectQueryResult struct { @@ -170,7 +185,7 @@ func (qt *IntersectQueryTree) getLocation(d *db.Database) (*db.Location, error) /////////////////////////////////////////////////////////////////////////////////////////////////// type CatQueryStep struct { *BaseQueryStep - Inputs []*uuid.UUID + Inputs []*util.GUID N int } @@ -250,10 +265,11 @@ type SetQueryTree struct { // Uses consistent hashing function to select node containing data for GET operation func (qt *SetQueryTree) getLocation(d *db.Database) (*db.Location, error) { slice, err := d.GetSliceForProfile(qt.profile_id) - fragment, err := d.GetFragmentForBitmap(slice, qt.bitmap) if err != nil { - return nil, err + //I should check GetFragmentForBitmap for possible errors but for now i'll just hardcode + return nil, NewFragmentNotFound(d.Name, qt.bitmap.FrameType, db.GetSlice(qt.profile_id)) } + fragment, err := d.GetFragmentForBitmap(slice, qt.bitmap) return fragment.GetLocation(), nil } @@ -379,17 +395,17 @@ func (qp *QueryPlanner) buildTree(query *Query, slice int) (QueryTree, error) { } // Produces flattened QueryPlan from QueryTree input -func (qp *QueryPlanner) flatten(qt QueryTree, id *uuid.UUID, location *db.Location) (*QueryPlan, error) { +func (qp *QueryPlanner) flatten(qt QueryTree, id *util.GUID, location *db.Location) (*QueryPlan, error) { plan := QueryPlan{} if cat, ok := qt.(*CatQueryTree); ok { - inputs := make([]*uuid.UUID, len(cat.subqueries)) + inputs := make([]*util.GUID, len(cat.subqueries)) loc, err := cat.getLocation(qp.Database) if err != nil { return nil, err } step := CatQueryStep{&BaseQueryStep{id, "cat", loc, location}, inputs, cat.N} for index, subq := range cat.subqueries { - sub_id := uuid.RandomUUID() + sub_id := util.RandomUUID() step.Inputs[index] = &sub_id subq_steps, err := qp.flatten(subq, &sub_id, loc) if err != nil { @@ -399,14 +415,14 @@ func (qp *QueryPlanner) flatten(qt QueryTree, id *uuid.UUID, location *db.Locati } plan = append(plan, step) } else if union, ok := qt.(*UnionQueryTree); ok { - inputs := make([]*uuid.UUID, len(union.subqueries)) + inputs := make([]*util.GUID, len(union.subqueries)) loc, err := union.getLocation(qp.Database) if err != nil { return nil, err } step := UnionQueryStep{&BaseQueryStep{id, "union", loc, location}, inputs} for index, subq := range union.subqueries { - sub_id := uuid.RandomUUID() + sub_id := util.RandomUUID() step.Inputs[index] = &sub_id subq_steps, err := qp.flatten(subq, &sub_id, loc) if err != nil { @@ -416,14 +432,14 @@ func (qp *QueryPlanner) flatten(qt QueryTree, id *uuid.UUID, location *db.Locati } plan = append(plan, step) } else if intersect, ok := qt.(*IntersectQueryTree); ok { - inputs := make([]*uuid.UUID, len(intersect.subqueries)) + inputs := make([]*util.GUID, len(intersect.subqueries)) loc, err := intersect.getLocation(qp.Database) if err != nil { return nil, err } step := IntersectQueryStep{&BaseQueryStep{id, "intersect", loc, location}, inputs} for index, subq := range intersect.subqueries { - sub_id := uuid.RandomUUID() + sub_id := util.RandomUUID() step.Inputs[index] = &sub_id subq_steps, err := qp.flatten(subq, &sub_id, loc) if err != nil { @@ -449,7 +465,7 @@ func (qp *QueryPlanner) flatten(qt QueryTree, id *uuid.UUID, location *db.Locati plan := QueryPlan{step} return &plan, nil } else if cnt, ok := qt.(*CountQueryTree); ok { - sub_id := uuid.RandomUUID() + sub_id := util.RandomUUID() loc, err := cnt.getLocation(qp.Database) if err != nil { return nil, err @@ -462,7 +478,7 @@ func (qp *QueryPlanner) flatten(qt QueryTree, id *uuid.UUID, location *db.Locati plan = append(plan, *subq_steps...) plan = append(plan, step) } else if topn, ok := qt.(*TopNQueryTree); ok { - sub_id := uuid.RandomUUID() + sub_id := util.RandomUUID() loc, err := topn.getLocation(qp.Database) if err != nil { return nil, err @@ -479,7 +495,7 @@ func (qp *QueryPlanner) flatten(qt QueryTree, id *uuid.UUID, location *db.Locati } // Transforms Query into QueryTree and flattens to QueryPlan object -func (qp *QueryPlanner) Plan(query *Query, id *uuid.UUID, destination *db.Location) (*QueryPlan, error) { +func (qp *QueryPlanner) Plan(query *Query, id *util.GUID, destination *db.Location) (*QueryPlan, error) { queryTree, err := qp.buildTree(query, -1) if err != nil { return nil, err diff --git a/query/planner_test.go b/query/planner_test.go index e1b4d9fcb..8980a1679 100644 --- a/query/planner_test.go +++ b/query/planner_test.go @@ -18,7 +18,7 @@ func basic_database() (*db.Database, *db.Fragment) { slice1 := database.GetOrCreateSlice(0) fragment_id1 := util.Id() fragment1 := database.GetOrCreateFragment(frame, slice1, fragment_id1) - process_id1 := uuid.RandomUUID() + process_id1 := util.RandomUUID() process1 := db.NewProcess(&process_id1) process1.SetHost("----192.1.1.0----") fragment1.SetProcess(process1) @@ -26,7 +26,7 @@ func basic_database() (*db.Database, *db.Fragment) { slice2 := database.GetOrCreateSlice(1) fragment_id2 := util.Id() fragment2 := database.GetOrCreateFragment(frame, slice2, fragment_id2) - process_id2 := uuid.RandomUUID() + process_id2 := util.RandomUUID() process2 := db.NewProcess(&process_id2) process2.SetHost("----192.1.1.1----") fragment2.SetProcess(process2) @@ -36,13 +36,13 @@ func basic_database() (*db.Database, *db.Fragment) { func TestQueryPlanner(t *testing.T) { Convey("Union query plan", t, func() { - id1 := uuid.RandomUUID() + id1 := util.RandomUUID() query1 := Query{Id: &id1, Operation: "get", Args: map[string]interface{}{"id": uint64(10), "frame": "general"}} - id2 := uuid.RandomUUID() + id2 := util.RandomUUID() query2 := Query{Id: &id2, Operation: "get", Args: map[string]interface{}{"id": uint64(20), "frame": "general"}} - id3 := uuid.RandomUUID() + id3 := util.RandomUUID() query := Query{Id: &id3, Operation: "union", Subqueries: []Query{query1, query2}} database, fragment1 := basic_database() @@ -50,7 +50,7 @@ func TestQueryPlanner(t *testing.T) { qplanner := QueryPlanner{Database: database, Query: &query} destination := fragment1.GetLocation() - id := uuid.RandomUUID() + id := util.RandomUUID() qpp, err := qplanner.Plan(&query, &id, destination) qp := *qpp @@ -63,7 +63,7 @@ func TestQueryPlanner(t *testing.T) { So(qp[1].(GetQueryStep).Slice, ShouldEqual, 0) So(*(qp[1].(GetQueryStep).Bitmap), ShouldResemble, db.Bitmap{20, "general", 0}) So(qp[2].(UnionQueryStep).Operation, ShouldEqual, "union") - So(qp[2].(UnionQueryStep).Inputs, ShouldResemble, []*uuid.UUID{ + So(qp[2].(UnionQueryStep).Inputs, ShouldResemble, []*GUID{ qp[0].(GetQueryStep).Id, qp[1].(GetQueryStep).Id, }) @@ -74,12 +74,12 @@ func TestQueryPlanner(t *testing.T) { So(qp[4].(GetQueryStep).Slice, ShouldEqual, 1) So(*(qp[4].(GetQueryStep).Bitmap), ShouldResemble, db.Bitmap{20, "general", 0}) So(qp[5].(UnionQueryStep).Operation, ShouldEqual, "union") - So(qp[5].(UnionQueryStep).Inputs, ShouldResemble, []*uuid.UUID{ + So(qp[5].(UnionQueryStep).Inputs, ShouldResemble, []*GUID{ qp[3].(GetQueryStep).Id, qp[4].(GetQueryStep).Id, }) So(qp[6].(CatQueryStep).Operation, ShouldEqual, "cat") - So(qp[6].(CatQueryStep).Inputs, ShouldResemble, []*uuid.UUID{ + So(qp[6].(CatQueryStep).Inputs, ShouldResemble, []*GUID{ qp[2].(UnionQueryStep).Id, qp[5].(UnionQueryStep).Id, }) @@ -94,7 +94,7 @@ func TestQueryPlanner(t *testing.T) { qplanner := QueryPlanner{Database: database, Query: query} destination := fragment1.GetLocation() - id := uuid.RandomUUID() + id := util.RandomUUID() qpp, err := qplanner.Plan(query, &id, destination) qp := *qpp So(err, ShouldEqual, nil) @@ -107,7 +107,7 @@ func TestQueryPlanner(t *testing.T) { So(qp[1].(GetQueryStep).Slice, ShouldEqual, 1) So(*(qp[1].(GetQueryStep).Bitmap), ShouldResemble, db.Bitmap{10, "general", 0}) So(qp[2].(CatQueryStep).Operation, ShouldEqual, "cat") - So(qp[2].(CatQueryStep).Inputs, ShouldResemble, []*uuid.UUID{ + So(qp[2].(CatQueryStep).Inputs, ShouldResemble, []*GUID{ qp[0].(GetQueryStep).Id, qp[1].(GetQueryStep).Id, }) @@ -122,7 +122,7 @@ func TestQueryPlanner(t *testing.T) { qplanner := QueryPlanner{Database: database, Query: query} destination := fragment1.GetLocation() - id := uuid.RandomUUID() + id := util.RandomUUID() qpp, err := qplanner.Plan(query, &id, destination) qp := *qpp So(err, ShouldEqual, nil) @@ -135,7 +135,7 @@ func TestQueryPlanner(t *testing.T) { So(qp[1].(GetQueryStep).Slice, ShouldEqual, 0) So(*(qp[1].(GetQueryStep).Bitmap), ShouldResemble, db.Bitmap{20, "general", 0}) So(qp[2].(UnionQueryStep).Operation, ShouldEqual, "union") - So(qp[2].(UnionQueryStep).Inputs, ShouldResemble, []*uuid.UUID{ + So(qp[2].(UnionQueryStep).Inputs, ShouldResemble, []*GUID{ qp[0].(GetQueryStep).Id, qp[1].(GetQueryStep).Id, }) @@ -146,12 +146,12 @@ func TestQueryPlanner(t *testing.T) { So(qp[4].(GetQueryStep).Slice, ShouldEqual, 1) So(*(qp[4].(GetQueryStep).Bitmap), ShouldResemble, db.Bitmap{20, "general", 0}) So(qp[5].(UnionQueryStep).Operation, ShouldEqual, "union") - So(qp[5].(UnionQueryStep).Inputs, ShouldResemble, []*uuid.UUID{ + So(qp[5].(UnionQueryStep).Inputs, ShouldResemble, []*GUID{ qp[3].(GetQueryStep).Id, qp[4].(GetQueryStep).Id, }) So(qp[6].(CatQueryStep).Operation, ShouldEqual, "cat") - So(qp[6].(CatQueryStep).Inputs, ShouldResemble, []*uuid.UUID{ + So(qp[6].(CatQueryStep).Inputs, ShouldResemble, []*GUID{ qp[2].(UnionQueryStep).Id, qp[5].(UnionQueryStep).Id, }) @@ -165,7 +165,7 @@ func TestQueryPlanner(t *testing.T) { qplanner := QueryPlanner{Database: database, Query: query} destination := fragment1.GetLocation() - id := uuid.RandomUUID() + id := util.RandomUUID() qpp, err := qplanner.Plan(query, &id, destination) qp := *qpp So(err, ShouldEqual, nil) @@ -183,7 +183,7 @@ func TestQueryPlanner(t *testing.T) { qplanner := QueryPlanner{Database: database, Query: query} destination := fragment1.GetLocation() - id := uuid.RandomUUID() + id := util.RandomUUID() qpp, err := qplanner.Plan(query, &id, destination) qp := *qpp So(err, ShouldEqual, nil) diff --git a/query/query.go b/query/query.go index 46748c94e..c3ba650c1 100644 --- a/query/query.go +++ b/query/query.go @@ -2,9 +2,8 @@ package query import ( "pilosa/db" + "pilosa/util" "strings" - - "github.com/gocql/gocql/uuid" ) type QueryInput interface{} @@ -16,13 +15,13 @@ type QueryResults struct { type PqlList []PqlListItem type PqlListItem struct { - Id *uuid.UUID + Id *util.GUID Label string PQL string } type Query struct { - Id *uuid.UUID + Id *util.GUID Operation string Args map[string]interface{} Subqueries []Query @@ -62,7 +61,7 @@ func QueryPlanForTokens(database *db.Database, tokens []Token, destination *db.L func QueryPlanForQuery(database *db.Database, query *Query, destination *db.Location) (*QueryPlan, error) { query_planner := QueryPlanner{Database: database, Query: query} - id := uuid.RandomUUID() + id := util.RandomUUID() query_plan, err := query_planner.Plan(query, &id, destination) if err != nil { return nil, err diff --git a/transport/tcp.go b/transport/tcp.go index 1a13562d7..09cfc314b 100644 --- a/transport/tcp.go +++ b/transport/tcp.go @@ -8,11 +8,11 @@ import ( "pilosa/config" "pilosa/core" "pilosa/db" + . "pilosa/util" "time" notify "github.com/bitly/go-notify" - "github.com/gocql/gocql/uuid" ) type connection struct { @@ -20,16 +20,16 @@ type connection struct { inbox chan *db.Message outbox chan *db.Message conn *net.Conn - process *uuid.UUID + process *GUID } type newconnection struct { - id *uuid.UUID + id *GUID connection *connection } func init() { - gob.Register(uuid.UUID{}) + gob.Register(GUID{}) } func (self *connection) manage() { @@ -78,7 +78,7 @@ BeginManageConnection: return } case message := <-self.inbox: - identifier, ok := message.Data.(uuid.UUID) + identifier, ok := message.Data.(GUID) if ok { // message is connection registration; bypass inbox and register self.process = &identifier @@ -103,7 +103,7 @@ type TcpTransport struct { port int inbox chan *db.Message outbox chan db.Envelope - connections map[uuid.UUID]*connection + connections map[GUID]*connection reg chan *newconnection } @@ -152,7 +152,7 @@ func (self *TcpTransport) Close() { log.Println("Shutting down TCP transport") } -func (self *TcpTransport) Send(message *db.Message, host *uuid.UUID) { +func (self *TcpTransport) Send(message *db.Message, host *GUID) { envelope := db.Envelope{message, host} notify.Post("outbox", &envelope) self.outbox <- envelope @@ -169,5 +169,5 @@ func (self *TcpTransport) Push(message *db.Message) { } func NewTcpTransport(service *core.Service) *TcpTransport { - return &TcpTransport{service, config.GetInt("port_tcp"), make(chan *db.Message, 100), make(chan db.Envelope, 100), make(map[uuid.UUID]*connection), make(chan *newconnection)} + return &TcpTransport{service, config.GetInt("port_tcp"), make(chan *db.Message, 100), make(chan db.Envelope, 100), make(map[GUID]*connection), make(chan *newconnection)} } diff --git a/util/id.go b/util/id.go index 7d001c09c..88b69fa0c 100644 --- a/util/id.go +++ b/util/id.go @@ -4,17 +4,28 @@ import ( "bytes" "encoding/binary" "encoding/hex" + "fmt" + "log" "math/rand" + "os" "strings" "time" + + "github.com/gocql/gocql" ) var ( counter = uint64(0) + Random *os.File ) func init() { rand.Seed(time.Now().UTC().UnixNano()) + f, err := os.Open("/dev/urandom") + if err != nil { + log.Fatal(err) + } + Random = f } type SUUID uint64 @@ -49,3 +60,51 @@ func Hex_to_SUUID(str string) SUUID { num := binary.BigEndian.Uint64(b) return SUUID(num) } + +type GUID [16]byte + +func (self GUID) String() string { + var offsets = [...]int{0, 2, 4, 6, 9, 11, 14, 16, 19, 21, 24, 26, 28, 30, 32, 34} + const hexString = "0123456789abcdef" + r := make([]byte, 36) + for i, b := range self { + r[offsets[i]] = hexString[b>>4] + r[offsets[i]+1] = hexString[b&0xF] + } + r[8] = '-' + r[13] = '-' + r[18] = '-' + r[23] = '-' + return string(r) + +} + +func RandomUUID() GUID { + uid, _ := gocql.RandomUUID() + var r GUID + copy(r[:], uid[:]) + return r +} +func ParseGUID(input string) (GUID, error) { + var u GUID + j := 0 + for _, r := range input { + switch { + case r == '-' && j&1 == 0: + continue + case r >= '0' && r <= '9' && j < 32: + u[j/2] |= byte(r-'0') << uint(4-j&1*4) + case r >= 'a' && r <= 'f' && j < 32: + u[j/2] |= byte(r-'a'+10) << uint(4-j&1*4) + case r >= 'A' && r <= 'F' && j < 32: + u[j/2] |= byte(r-'A'+10) << uint(4-j&1*4) + default: + return GUID{}, fmt.Errorf("invalid GUID %q", input) + } + j += 1 + } + if j != 32 { + return GUID{}, fmt.Errorf("invalid GUID %q", input) + } + return u, nil +} diff --git a/util/util_test.go b/util/util_test.go index e90974734..7fdfa4fcc 100644 --- a/util/util_test.go +++ b/util/util_test.go @@ -1,9 +1,10 @@ package util import ( + "fmt" "testing" - "github.com/gocql/gocql/uuid" + "github.com/gocql/gocql" . "github.com/smartystreets/goconvey/convey" ) @@ -11,14 +12,14 @@ import ( var ( array [1000000]int muid = make(map[SUUID]int) - muuid = make(map[*uuid.UUID]int) + muuid = make(map[*GUID]int) r int ) func init() { for i, _ := range array { muid[Id()] = i - id := uuid.RandomUUID() + id := util.RandomUUID() muuid[&id] = i } @@ -66,9 +67,13 @@ func BenchmarkId(b *testing.B) { func BenchmarkUUID(b *testing.B) { // run the Fib function b.N times for n := 0; n < b.N; n++ { - uuid.RandomUUID() + gocql.RandomUUID() } } +func TestGUID(t *testing.T) { + fmt.Println(RandomUUID().String()) + +} /* func BenchmarkLookupId(b *testing.B) { @@ -80,7 +85,7 @@ func BenchmarkLookupId(b *testing.B) { } } func BenchmarkLookupUUID(b *testing.B) { - x := uuid.RandomUUID() + x := util.RandomUUID() for i := 0; i < b.N; i++ { if a, found := muuid[&x]; found { muuid[&x] = a + 1