diff --git a/core/http.go b/core/http.go index 85194a3cb..0cd08ab2f 100644 --- a/core/http.go +++ b/core/http.go @@ -436,7 +436,7 @@ func (self *WebService) HandleSetBit(w http.ResponseWriter, r *http.Request) { http.Error(w, "Missing db", http.StatusBadRequest) return } - db := obj["db"].(string) + dbs := obj["db"].(string) if obj["frame"] == nil { http.Error(w, "Missing Frame", http.StatusBadRequest) @@ -451,16 +451,33 @@ func (self *WebService) HandleSetBit(w http.ResponseWriter, r *http.Request) { } t = float64(obj["filter"].(float64)) filter := int(t) - + result := false for bitmap_id := range bitmaps(frame, obj) { - pql := fmt.Sprintf("set(%d, %s, %d, %d)", bitmap_id, frame, filter, profile_id) start := time.Now() - result, err := self.service.Executor.RunPQL(db, pql) + + database := self.service.Cluster.GetOrCreateDatabase(dbs) + frag, err := database.GetFragmentFromProfile(frame, profile_id) + if err != nil { + //no fragment + self.service.TopologyMapper.MakeFragments(dbs, db.GetSlice(profile_id)) + time.Sleep(2 * time.Second) + continue + } + if util.Equal(frag.GetProcessId(), self.service.Id) { + // The Local Route + result, _ = self.service.Index.SetBit(frag.GetId(), bitmap_id, profile_id, uint64(filter)) + } else { + println("remote") + // The Remote Route + } + + //result, err := self.service.Executor.RunPQL(db, pql) + //pql := fmt.Sprintf("set(%d, %s, %d, %d)", bitmap_id, frame, filter, profile_id) delta := time.Since(start) util.SendTimer("executor_setbit", delta.Nanoseconds()) if err != nil { - log.Println("Error running set_bit", db, pql) + log.Println("Error running set_bit", dbs, frame, profile_id) http.Error(w, err.Error(), http.StatusInternalServerError) return } @@ -470,7 +487,6 @@ func (self *WebService) HandleSetBit(w http.ResponseWriter, r *http.Request) { } } - encoder := json.NewEncoder(w) err = encoder.Encode(results) if err != nil { diff --git a/db/topology.go b/db/topology.go index c71d5ef7f..b706fc6e9 100644 --- a/db/topology.go +++ b/db/topology.go @@ -387,6 +387,16 @@ func (d *Database) GetFragmentById(fragment_id *GUID) *Fragment { } */ +func (d *Database) GetFragmentFromProfile(frame string, profile_id uint64) (*Fragment, error) { + slice := GetSlice(profile_id) + for _, frameslice := range d.frame_slice_intersects { + if frameslice.frame.name == frame && frameslice.slice.id == slice { + return frameslice.fragment, nil + } + } + return nil, errors.New("FragmentDoes does not exist.") +} + func (d *Database) getFragment(frame *Frame, slice *Slice) (*Fragment, error) { fsi, err := d.GetFrameSliceIntersect(frame, slice) if err != nil { diff --git a/index/brand.go b/index/brand.go index c777f0c23..2696b79ee 100644 --- a/index/brand.go +++ b/index/brand.go @@ -127,10 +127,10 @@ func (self *Brand) SetBit(bitmap_id uint64, bit_pos uint64, filter uint64) bool change, chunk, address := SetBit(bm, bit_pos) if change { val := chunk.Value.Block[address.BlockIndex] - /*self.storage.BeginBatch() - self.storage.StoreBlock(bitmap_id, self.db, self.frame, self.slice, filter, address.ChunkKey, int32(address.BlockIndex), val) - self.storage.StoreBlock(bitmap_id, self.db, self.frame, self.slice, filter, COUNTERMASK, 0, bm.Count()) - self.storage.EndBatch() + /* self.storage.BeginBatch() + self.storage.StoreBlock(bitmap_id, self.db, self.frame, self.slice, filter, address.ChunkKey, int32(address.BlockIndex), val) + self.storage.StoreBlock(bitmap_id, self.db, self.frame, self.slice, filter, COUNTERMASK, 0, bm.Count()) + self.storage.EndBatch() */ self.storage.StoreBit(bitmap_id, self.db, self.frame, self.slice, filter, address.ChunkKey, int32(address.BlockIndex), val, bm.Count()) self.rank_count++ diff --git a/index/storage_cass.go b/index/storage_cass.go index 28e19f642..c0b403c0a 100644 --- a/index/storage_cass.go +++ b/index/storage_cass.go @@ -213,7 +213,8 @@ type CassQueue struct { } func NewCassQueue() CassQueue { - return CassQueue{0, make(chan CassRecord, 2048)} + //return CassQueue{0, make(chan CassRecord, 4096)} + return CassQueue{0, make(chan CassRecord)} } func (self *CassQueue) Push(rec CassRecord) {