diff --git a/db/topology.go b/db/topology.go index 52c99edb7..08f0bbf47 100644 --- a/db/topology.go +++ b/db/topology.go @@ -16,19 +16,11 @@ var SliceDoesNotExistError = errors.New("Slice does not exist.") var FragmentDoesNotExistError = errors.New("Fragment does not exist.") var FrameSliceIntersectDoesNotExistError = errors.New("FrameSliceIntersect does not exist.") -/* -type Location struct { - Ip string - Port int -} -*/ type Location struct { ProcessId *uuid.UUID FragmentId util.SUUID } -type ProcessId *uuid.UUID - type Process struct { id *uuid.UUID host string diff --git a/dispatch/dispatch.go b/dispatch/dispatch.go index cf55b0ead..df48a07a3 100644 --- a/dispatch/dispatch.go +++ b/dispatch/dispatch.go @@ -32,7 +32,7 @@ func (self *Dispatch) Run() { switch message.Data.(type) { case query.GetQueryStep, query.SetQueryStep: fmt.Println("GET/SET QUERYSTEP") - self.service.Executor.NewJob(message) + go self.service.Executor.NewJob(message) default: fmt.Println("unknown") } diff --git a/executor/executor.go b/executor/executor.go index 9bf992480..9e026ace2 100644 --- a/executor/executor.go +++ b/executor/executor.go @@ -23,7 +23,6 @@ type Executor struct { service *core.Service inbox chan *db.Message qs_chan chan *query.QueryStep - hold map[*uuid.UUID]chan *query.QueryResults } func (self *Executor) Init() error { @@ -36,11 +35,8 @@ func (self *Executor) Close() { } func (self *Executor) NewJob(job *db.Message) { - //j := Job{qp, results} - //self.inbox <- &j spew.Dump("NewJob") spew.Dump(job.Data) - // TODO: switch on job.Data type switch job.Data.(type) { case query.GetQueryStep: @@ -65,8 +61,9 @@ func (self *Executor) NewJob(job *db.Message) { spew.Dump(err) } spew.Dump(count) + spew.Dump(qs.Id) - //self.Set(job.Data.Id, query_results) + //self.Set(qs.Id, count) case query.SetQueryStep: fmt.Println("SET QUERYSTEP") //self.Get() @@ -76,13 +73,6 @@ func (self *Executor) NewJob(job *db.Message) { } } -/* -func (self *Executor) NewJob(qp *query.QueryPlan, results chan *query.QueryResults) { - j := Job{qp, results} - self.inbox <- &j -} -*/ - func (self *Executor) NewQS(qs *query.QueryStep) { self.qs_chan <- qs } @@ -168,5 +158,5 @@ func (self *Executor) Run() { } func NewExecutor(service *core.Service) *Executor { - return &Executor{service, make(chan *db.Message), make(chan *query.QueryStep), make(map[*uuid.UUID]chan *query.QueryResults)} + return &Executor{service, make(chan *db.Message), make(chan *query.QueryStep)} }