auto allocate .001

This commit is contained in:
Todd Gruben 2014-06-18 15:09:31 -05:00
parent 4a2c3e73d0
commit 1b89d79c0a
20 changed files with 250 additions and 140 deletions

View file

@ -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)

View file

@ -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)
}

View file

@ -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

View file

@ -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)

View file

@ -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)
}

View file

@ -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{}
}

View file

@ -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 {
}
*/

View file

@ -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
}

View file

@ -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()
}
*/

View file

@ -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)
}()

View file

@ -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

View file

@ -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()

View file

@ -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)
}

View file

@ -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()

View file

@ -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

View file

@ -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)

View file

@ -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

View file

@ -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)}
}

View file

@ -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
}

View file

@ -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