diff --git a/index/commands.go b/index/commands.go index 7447c58a7..18b2d5fb9 100644 --- a/index/commands.go +++ b/index/commands.go @@ -41,6 +41,7 @@ type Command interface { Execute(*Fragment) Calculation GetResponder() *Responder } + type CmdGet struct { meta *Responder bitmap_id uint64 @@ -125,3 +126,38 @@ func (cmd *CmdSetBit) Execute(f *Fragment) Calculation { func (cmd *CmdSetBit) GetResponder() *Responder { return cmd.meta } + +type CmdGetBytes struct { + meta *Responder + bitmap BitmapHandle +} + +func NewGetBytes(bh BitmapHandle) *CmdGetBytes { + return &CmdGetBytes{NewResponder("GetBytes"), bh} +} + +func (cmd *CmdGetBytes) GetResponder() *Responder { + return cmd.meta +} +func (cmd *CmdGetBytes) Execute(f *Fragment) Calculation { + bm, _ := f.getBitmap(cmd.bitmap) + return bm.ToBytes() +} + +type CmdFromBytes struct { + meta *Responder + bytes []byte +} + +func NewFromBytes(bytes []byte) *CmdFromBytes { + return &CmdFromBytes{NewResponder("FromBytes"), bytes} +} + +func (cmd *CmdFromBytes) GetResponder() *Responder { + return cmd.meta +} +func (cmd *CmdFromBytes) Execute(f *Fragment) Calculation { + result := NewBitmap() + result.FromBytes(cmd.bytes) + return f.AllocHandle(result) +} diff --git a/index/server.go b/index/server.go index ef4d91a94..e6d0f71da 100644 --- a/index/server.go +++ b/index/server.go @@ -58,6 +58,24 @@ func (self *FragmentContainer) Count(frag_id SUUID, bitmap BitmapHandle) (uint64 return 0, errors.New("Invalid Bitmap Handle") } +func (self *FragmentContainer) GetBytes(frag_id SUUID, bh BitmapHandle) ([]byte, error) { + if fragment, found := self.GetFragment(frag_id); found { + request := NewGetBytes(bh) + fragment.requestChan <- request + return request.GetResponder().Response().answer.([]byte), nil + } + return nil, errors.New("Invalid Bitmap Handle") +} + +func (self *FragmentContainer) FromBytes(frag_id SUUID, bytes []byte) (BitmapHandle, error) { + if fragment, found := self.GetFragment(frag_id); found { + request := NewFromBytes(bytes) + fragment.requestChan <- request + return request.GetResponder().Response().answer.(BitmapHandle), nil + } + return 0, errors.New("Invalid Bitmap Handle") +} + func (self *FragmentContainer) SetBit(frag_id SUUID, bitmap_id uint64, pos uint64) (bool, error) { if fragment, found := self.GetFragment(frag_id); found { request := NewSetBit(bitmap_id, pos) diff --git a/index/server_test.go b/index/server_test.go index eb930ec36..22af9f8d0 100644 --- a/index/server_test.go +++ b/index/server_test.go @@ -52,4 +52,15 @@ func TestServer(t *testing.T) { num, _ = dummy.Count(id, result) So(num, ShouldEqual, 0) }) + Convey("Bytes", t, func() { + bi1 := uint64(1234) + bh1, _ := dummy.Get(id, bi1) + before, _ := dummy.Count(id, bh1) + + bytes, _ := dummy.GetBytes(id, bh1) + bh2, _ := dummy.FromBytes(id, bytes) + + after, _ := dummy.Count(id, bh2) + So(before, ShouldEqual, after) + }) } diff --git a/index/storage_kvlite.go b/index/storage_kvlite.go new file mode 100644 index 000000000..6b0ce3642 --- /dev/null +++ b/index/storage_kvlite.go @@ -0,0 +1,124 @@ +package index + +// #cgo CFLAGS:-mpopcnt + +import ( + "bytes" + "encoding/binary" + "fmt" + "os" + + "github.com/steveyen/gkvlite" +) + +type KVStorage struct { + db *gkvlite.Store + cc *gkvlite.Collection +} + +func NewKVStorage(path string, slice int, db string) (Storage, error) { + // log.Println("Hello") + obj := new(KVStorage) + + f, err := os.Create(fmt.Sprintf("%sKV.%s.%d", path, db, slice)) + s, err := gkvlite.NewStore(f) + obj.cc = s.SetCollection(db, nil) + + obj.db = s + + return obj, err +} + +func (self *KVStorage) Fetch(bitmap_id uint64, db string, slice int) IBitmap { + + var ( + chunk *Chunk + ) + last_key := COUNTER_KEY + marker := COUNTER_KEY + count := uint64(0) + + key, _ := toKeyBytes(int64(bitmap_id), slice, 0, 0) + + bitmap := CreateRBBitmap() + + self.cc.VisitItemsAscend(key, true, func(i *gkvlite.Item) bool { + bid, _, chunk_key, block_index, _ := fromKeyBytes(i.Key) + if bid != int64(bitmap_id) { + return false + } + block, _ := binary.Uvarint(i.Val) //just need to cast as a block + if chunk_key != marker { + if chunk_key != last_key { + chunk = &Chunk{uint64(chunk_key), BlockArray{}} + bitmap.AddChunk(chunk) + } + chunk.Value.Block[block_index] = block + + } else { + count = block + } + last_key = chunk_key + return true + }) + bitmap.SetCount(uint64(count)) + return bitmap +} + +func (self *KVStorage) Store(bitmap_id int64, db string, slice int, bitmap *Bitmap) error { + for i := bitmap.Min(); !i.Limit(); i = i.Next() { + var chunk = i.Item() + for idx, block := range chunk.Value.Block { + block_index := int32(idx) + iblock := int64(block) + if iblock != 0 { + self.StoreBlock(bitmap_id, db, slice, int64(chunk.Key), block_index, iblock) + } + } + } + cnt := int64(BitCount(bitmap)) + + //var dumb = COUNTERMASK + //COUNTER_KEY := int64(dumb) + + self.StoreBlock(bitmap_id, db, slice, COUNTER_KEY, 0, cnt) + self.db.Flush() + return nil + +} +func Bytes(src int64) []byte { + buf := new(bytes.Buffer) + err := binary.Write(buf, binary.LittleEndian, src) + if err != nil { + fmt.Println("binary.Write failed:", err) + } + return buf.Bytes() +} +func fromKeyBytes(src []byte) (bitmap_id int64, slice int, chunk_key int64, block_index int32, err error) { + buf := bytes.NewReader(src) + err = binary.Read(buf, binary.LittleEndian, &bitmap_id) + err = binary.Read(buf, binary.LittleEndian, &slice) + err = binary.Read(buf, binary.LittleEndian, &chunk_key) + err = binary.Read(buf, binary.LittleEndian, &block_index) + return +} + +func toKeyBytes(bitmap_id int64, slice int, chunk_key int64, block_index int32) ([]byte, error) { + buf := new(bytes.Buffer) + err := binary.Write(buf, binary.LittleEndian, bitmap_id) + err = binary.Write(buf, binary.LittleEndian, slice) + err = binary.Write(buf, binary.LittleEndian, chunk_key) + err = binary.Write(buf, binary.LittleEndian, block_index) + + return buf.Bytes(), err +} + +func (self *KVStorage) StoreBlock(bitmap_id int64, db string, slice int, chunk_key int64, block_index int32, block int64) error { + //key := fmt.Sprintf("%d%d:%d:%d", bitmap_id, slice, chunk_key, block_index) + key, _ := toKeyBytes(bitmap_id, slice, chunk_key, block_index) + value := Bytes(block) + + self.cc.Set(key, value) + + return nil +} diff --git a/index/storage_test.go b/index/storage_test.go new file mode 100644 index 000000000..c998324e9 --- /dev/null +++ b/index/storage_test.go @@ -0,0 +1,29 @@ +package index + +import ( + "testing" + + // "io/ioutil" + // "time" + . "github.com/smartystreets/goconvey/convey" +) + +func TestStorage(t *testing.T) { + db := "db" + slice := 0 + bitmap_id := uint64(1234) + Convey("KV ", t, func() { + storage, _ := NewKVStorage("/tmp/", 0, db) + bm := storage.Fetch(bitmap_id, db, slice) + SetBit(bm, 0) + SetBit(bm, 1) + SetBit(bm, 2) + storage.Store(int64(bitmap_id), db, slice, bm.(*Bitmap)) + bm2 := storage.Fetch(bitmap_id, db, slice) + So(BitCount(bm), ShouldEqual, BitCount(bm2)) + So(BitCount(bm), ShouldEqual, bm.Count()) + So(BitCount(bm), ShouldEqual, 3) + + }) + +}