added a temp HTTP listener to transport.tcp so i could test the dispatch calls to process (index)

This commit is contained in:
travisturner 2013-12-19 14:04:51 -06:00
parent 9ee028290b
commit 6d7c683a8a
8 changed files with 162 additions and 12 deletions

View file

@ -6,7 +6,8 @@ import (
"log"
"os"
"sync"
"github.com/nu7hatch/gouuid"
"github.com/nu7hatch/gouuid"
"launchpad.net/goyaml"
)

View file

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

View file

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

View file

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

View file

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

View file

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

View file

@ -5,7 +5,7 @@ import (
)
type Transporter interface {
Init() error
Run()
Close()
Send(*db.Message)
Receive() *db.Message

View file

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