diff --git a/config/config.go b/config/config.go index 16fb8de0f..b528b228a 100644 --- a/config/config.go +++ b/config/config.go @@ -6,7 +6,8 @@ import ( "log" "os" "sync" - "github.com/nu7hatch/gouuid" + + "github.com/nu7hatch/gouuid" "launchpad.net/goyaml" ) diff --git a/core/etcd.go b/core/etcd.go index 35a2f9e3c..2ebc91188 100644 --- a/core/etcd.go +++ b/core/etcd.go @@ -5,11 +5,11 @@ import ( "log" "pilosa/config" "pilosa/db" + "pilosa/util" "strconv" "strings" "sync" - "pilosa/util" "github.com/coreos/go-etcd/etcd" "github.com/davecgh/go-spew/spew" "github.com/nu7hatch/gouuid" @@ -29,7 +29,6 @@ func (self *TopologyMapper) Run() { for _, node := range flatten(resp.Node) { err := self.handlenode(node) if err != nil { - spew.Dump(node) log.Println(err) } } @@ -63,6 +62,7 @@ func (self *TopologyMapper) handlenode(node *etcd.Node) error { var fragment *db.Fragment var fragment_id util.SUUID var slice *db.Slice + var slice_int int var process_uuid *uuid.UUID var process *db.Process var err error @@ -87,7 +87,7 @@ func (self *TopologyMapper) handlenode(node *etcd.Node) error { } } if len(bits) > 5 { - slice_int, err := strconv.Atoi(bits[5]) + slice_int, err = strconv.Atoi(bits[5]) if err != nil { return err } @@ -117,8 +117,9 @@ func (self *TopologyMapper) handlenode(node *etcd.Node) error { process = db.NewProcess(process_uuid) fragment.SetProcess(process) - //if service.process_id == process_uuid: - // service.index.AddFragment(bits[3], bits[1], slice_int, fragment_uuid) + if self.service.process_id.String() == process_uuid.String() { + self.service.Process.AddFragment(bits[1], bits[3], slice_int, fragment_id) + } } return err diff --git a/core/service.go b/core/service.go index 65d4a7ffa..eb160e288 100644 --- a/core/service.go +++ b/core/service.go @@ -6,10 +6,12 @@ import ( "os/signal" "pilosa/config" "pilosa/db" + "pilosa/index" "pilosa/interfaces" "syscall" "github.com/coreos/go-etcd/etcd" + "github.com/davecgh/go-spew/spew" "github.com/nu7hatch/gouuid" ) @@ -25,9 +27,11 @@ type Service struct { Dispatch interfaces.Dispatcher WebService *WebService process_id *uuid.UUID + Process *index.FragmentContainer } func NewService() *Service { + spew.Dump("NewService") service := new(Service) service.init_id() service.Etcd = etcd.NewClient(nil) @@ -37,6 +41,7 @@ func NewService() *Service { service.ProcessMap = NewProcessMap() service.WebService = NewWebService(service) service.process_id = config.GetUUID("process_id") + service.Process = index.NewFragmentContainer() return service } @@ -72,6 +77,8 @@ func (service *Service) Run() { go service.TopologyMapper.Run() go service.ProcessMapper.Run() go service.WebService.Run() + go service.Transport.Run() + go service.Dispatch.Run() sigterm, sighup := service.GetSignals() for { diff --git a/cruncher/cruncher.go b/cruncher/cruncher.go index 919f3dc14..b5c17742e 100644 --- a/cruncher/cruncher.go +++ b/cruncher/cruncher.go @@ -1,11 +1,12 @@ package cruncher import ( - "github.com/davecgh/go-spew/spew" "pilosa/core" "pilosa/dispatch" "pilosa/index" "pilosa/transport" + + "github.com/davecgh/go-spew/spew" ) type Cruncher struct { diff --git a/dispatch/dispatch.go b/dispatch/dispatch.go index 90ba94b19..14a4ab305 100644 --- a/dispatch/dispatch.go +++ b/dispatch/dispatch.go @@ -3,6 +3,12 @@ package dispatch import ( "log" "pilosa/core" + "pilosa/index" + "pilosa/util" + "strconv" + "strings" + + "github.com/davecgh/go-spew/spew" ) type CruncherDispatch struct { @@ -19,9 +25,98 @@ func (self *CruncherDispatch) Close() { } func (self *CruncherDispatch) Run() { + log.Println("Dispatch Run...") for { message := self.service.Transport.Receive() log.Println("Processing ", message) + spew.Dump(message.Key) + spew.Dump(message.Data) + + path := message.Data.(string) + + bits := strings.Split(path, "/") + + var fragment_id util.SUUID + var bitmaps []uint64 + var profile_id uint64 + var s uint64 + + command := bits[0] + if len(bits) > 1 { + fragment_id = util.Hex_to_SUUID(bits[1]) + } + if len(bits) > 2 { + bitmap_ids := strings.Split(bits[2], ",") + spew.Dump(bitmap_ids) + for i := range bitmap_ids { + spew.Dump(i, bitmap_ids[i]) + s, _ = strconv.ParseUint(bitmap_ids[i], 10, 64) + bitmaps = append(bitmaps, s) + } + } + if len(bits) > 3 { + profile_id, _ = strconv.ParseUint(bits[3], 10, 64) + } + + spew.Dump("COMMAND:", command) + spew.Dump("FRAGID:", fragment_id) + spew.Dump("BITMAPS:", bitmaps) + spew.Dump("PROFILEID:", profile_id) + + if command == "set" { + res, err := self.service.Process.SetBit(fragment_id, bitmaps[0], profile_id) + spew.Dump("SET") + spew.Dump(res) + spew.Dump(err) + } + if command == "count" { + spew.Dump("COUNT") + bh, err := self.service.Process.Get(fragment_id, bitmaps[0]) + if err != nil { + spew.Dump(err) + } + count, err := self.service.Process.Count(fragment_id, bh) + if err != nil { + spew.Dump(err) + } + spew.Dump(count) + } + if command == "intersect" { + spew.Dump("INTERSECT") + var bhs []index.BitmapHandle + for i := range bitmaps { + bh, _ := self.service.Process.Get(fragment_id, bitmaps[i]) + bhs = append(bhs, bh) + } + bhi, err := self.service.Process.Intersect(fragment_id, bhs) + if err != nil { + spew.Dump(err) + } + + count, err := self.service.Process.Count(fragment_id, bhi) + if err != nil { + spew.Dump(err) + } + spew.Dump(count) + } + if command == "union" { + spew.Dump("UNION") + var bhs []index.BitmapHandle + for i := range bitmaps { + bh, _ := self.service.Process.Get(fragment_id, bitmaps[i]) + bhs = append(bhs, bh) + } + bhi, err := self.service.Process.Union(fragment_id, bhs) + if err != nil { + spew.Dump(err) + } + + count, err := self.service.Process.Count(fragment_id, bhi) + if err != nil { + spew.Dump(err) + } + spew.Dump(count) + } } } diff --git a/index/fragment_container.go b/index/fragment_container.go index 77e83a8f4..04c1049f9 100644 --- a/index/fragment_container.go +++ b/index/fragment_container.go @@ -115,7 +115,7 @@ func (self *FragmentContainer) SetBit(frag_id SUUID, bitmap_id uint64, pos uint6 return false, errors.New("Invalid Bitmap Handle") } -func (self *FragmentContainer) AddFragment(frame string, db string, slice int, id SUUID) { +func (self *FragmentContainer) AddFragment(db string, frame string, slice int, id SUUID) { f := NewFragment(id, db, slice, frame) self.fragments[id] = f go f.ServeFragment() diff --git a/interfaces/core.go b/interfaces/core.go index e6d991820..1143370bf 100644 --- a/interfaces/core.go +++ b/interfaces/core.go @@ -5,7 +5,7 @@ import ( ) type Transporter interface { - Init() error + Run() Close() Send(*db.Message) Receive() *db.Message diff --git a/transport/tcp.go b/transport/tcp.go index e30cbde09..f8143b6a5 100644 --- a/transport/tcp.go +++ b/transport/tcp.go @@ -1,19 +1,64 @@ package transport import ( + "fmt" "log" + "net" + "net/http" + "pilosa/config" "pilosa/core" "pilosa/db" + "time" ) type TcpTransport struct { port int inbox chan *db.Message + stop chan bool } -func (self *TcpTransport) Init() error { +func (self *TcpTransport) RunServer(porti int) { + http.Handle("/", self) + port := fmt.Sprintf(":%d", porti) + + s := &http.Server{ + Addr: port, + Handler: nil, + ReadTimeout: 10 * time.Second, + WriteTimeout: 10 * time.Second, + MaxHeaderBytes: 1 << 20, + } + + l, e := net.Listen("tcp", port) + if e != nil { + log.Panicf(e.Error()) + } + self.stop = make(chan bool) + go s.Serve(l) + select { + case <-self.stop: + log.Printf("Server thread exit") + l.Close() + // Shutdown() + return + break + } +} + +func (self *TcpTransport) ServeHTTP(w http.ResponseWriter, r *http.Request) { + //fmt.Fprintf(w, "URL Path: %s\n", r.URL.Path[1:]) + path := r.URL.Path[1:] + + msg := new(db.Message) + msg.Key = "path" + msg.Data = path + self.inbox <- msg +} + +func (self *TcpTransport) Run() { log.Println("Initializing TCP transport") - return nil + + self.RunServer(self.port) } func (self *TcpTransport) Close() { @@ -29,5 +74,5 @@ func (self *TcpTransport) Receive() *db.Message { } func NewTcpTransport(service *core.Service) *TcpTransport { - return &TcpTransport{12000, make(chan *db.Message)} //TODO: make port configurable + return &TcpTransport{config.GetInt("port_tcp"), make(chan *db.Message), nil} }