From 0cf2ddb08678277446f88067004c009bad235350 Mon Sep 17 00:00:00 2001 From: travisturner Date: Tue, 10 Dec 2013 14:26:59 -0600 Subject: [PATCH] first pass at GetOrCreate methods in db.topology --- core/etcd.go | 8 +- db/topology.go | 304 +++++++++++++++++++++++++++++++------------- db/topology_test.go | 25 +++- 3 files changed, 246 insertions(+), 91 deletions(-) diff --git a/core/etcd.go b/core/etcd.go index 95eaeb093..795ce4b1f 100644 --- a/core/etcd.go +++ b/core/etcd.go @@ -121,11 +121,13 @@ func (service *Service) MetaWatcher() { key = slice_attr_ref.Key[len(slice_ref.Key)+1:] if key == "fragment" { for _, fragment_ref := range slice_attr_ref.Nodes { + /* fragment_name := fragment_ref.Key[len(slice_attr_ref.Key)+1:] fragment_id, err := strconv.Atoi(fragment_name) if err != nil { log.Fatal(err) } + */ for _, fragment_attr_ref := range fragment_ref.Nodes { key = fragment_attr_ref.Key[len(fragment_ref.Key)+1:] if key == "node" { @@ -133,8 +135,8 @@ func (service *Service) MetaWatcher() { if err != nil { log.Fatal(err) } - process := db.NewProcess(uuid) - database.AddFragment(frame, slice, process, fragment_id) + process := uuid + database.AddFragment(frame, slice, process) } } } @@ -149,7 +151,7 @@ func (service *Service) MetaWatcher() { } database, _ := cluster.GetDatabase("main") spew.Dump(database) - database.TestSetBit(db.Bitmap{1200, "general"}, 1) + database.OldGetFragment(db.Bitmap{1200, "general"}, 1) receiver := make(chan *etcd.Response) stop := make(chan bool) diff --git a/db/topology.go b/db/topology.go index 15c61561b..226b48404 100644 --- a/db/topology.go +++ b/db/topology.go @@ -2,13 +2,14 @@ package db import ( "github.com/stathat/consistent" - "github.com/davecgh/go-spew/spew" + //"github.com/davecgh/go-spew/spew" "github.com/nu7hatch/gouuid" "log" "fmt" "errors" "strings" "strconv" + "sync" ) var FrameDoesNotExistError = errors.New("Frame does not exist.") @@ -50,52 +51,15 @@ func (location *Location) ToString() string { // Map of node location to their router type NodeMap map[Location]Location -// A fragment is a collection of bitmaps within a slice. The fragment contains a reference to the responsible node for that fragment. The node is in the form ip:port -type Fragment struct { - process *Process - id int -} -// A slice is the vertical combination of every fragment. It contains the hashring used to delegate bitmaps to fragments -type Slice struct { - id int -} -// A frame is a collection of slices in a given category (brands, demographics, etc), specific to a database -type Frame struct { - name string -} -type FrameSliceIntersect struct { - slice *Slice - frame *Frame - fragments []*Fragment - hashring *consistent.Consistent -} - -func (fsi *FrameSliceIntersect) GetFragment(fragment_id int) (*Fragment, error) { - for _, fragment := range fsi.fragments { - if fragment.id == fragment_id { - return fragment, nil - } - } - return nil, FragmentDoesNotExistError -} - -// Add a slice to a database -func (d *Database) AddSlice(slice_id int) *Slice { - slice := Slice{id: slice_id} - d.slices = append(d.slices, &slice) - // add intersections - for _, frame := range d.frames { - d.AddFrameSliceIntersect(frame, &slice) - } - return &slice -} +/////////// CLUSTERS //////////////////////////////////////////////////////////////////// // Represents the entire cluster, and a reference to the Node this instance is running on type Cluster struct { databases map[string]*Database + mutex sync.Mutex } func NewCluster() *Cluster { @@ -104,8 +68,22 @@ func NewCluster() *Cluster { return &cluster } + +/////////// DATABASES //////////////////////////////////////////////////////////////////// + +// A database is a collection of all the frames within a given profile space +type Database struct { + Name string + frames []*Frame + slices []*Slice + frame_slice_intersects []*FrameSliceIntersect + mutex sync.Mutex +} + // Add a database to a cluster func (c *Cluster) AddDatabase(name string) *Database { + c.mutex.Lock() + defer c.mutex.Unlock() database := Database{Name: name} if c.databases == nil { c.databases = make(map[string]*Database) @@ -115,6 +93,8 @@ func (c *Cluster) AddDatabase(name string) *Database { } func (c *Cluster) GetDatabase(name string) (*Database, error) { + c.mutex.Lock() + defer c.mutex.Unlock() value, ok := c.databases[name] if !ok { return nil, errors.New("The database does not exist!") @@ -123,12 +103,14 @@ func (c *Cluster) GetDatabase(name string) (*Database, error) { } } -// A database is a collection of all the frames within a given profile space -type Database struct { - Name string - frames []*Frame - slices []*Slice - frame_slice_intersects []*FrameSliceIntersect +func (c *Cluster) GetOrCreateDatabase(name string) *Database { + c.mutex.Lock() + defer c.mutex.Unlock() + database, err := c.GetDatabase(name) + if err == nil { + return database + } + return c.AddDatabase(name) } // Count the number of slices in a database @@ -139,8 +121,30 @@ func (d *Database) NumSlices() (int, error) { return len(d.slices), nil } +///////// FRAMES //////////////////////////////////////////////////////////////////// + +// A frame is a collection of slices in a given category +// (brands, demographics, etc), specific to a database +type Frame struct { + name string +} + +// Get a frame from a database +func (d *Database) GetFrame(name string) (*Frame, error) { + d.mutex.Lock() + defer d.mutex.Unlock() + for _, frame := range d.frames { + if frame.name == name { + return frame, nil + } + } + return nil, FrameDoesNotExistError +} + // Add a frame to a database func (d *Database) AddFrame(name string) *Frame { + d.mutex.Lock() + defer d.mutex.Unlock() frame := Frame{name: name} d.frames = append(d.frames, &frame) // add intersections @@ -150,24 +154,74 @@ func (d *Database) AddFrame(name string) *Frame { return &frame } -func (d *Database) AddFragment(frame *Frame, slice *Slice, process *Process, fragment_id int) *Fragment { +func (d *Database) GetOrCreateFrame(name string) *Frame { + d.mutex.Lock() + defer d.mutex.Unlock() + frame, err := d.GetFrame(name) + if err == nil { + return frame + } + return d.AddFrame(name) +} - frameslice, _ := d.GetFrameSliceIntersect(frame, slice) - fragment := Fragment{process: process, id: fragment_id} - frameslice.fragments = append(frameslice.fragments, &fragment) - frameslice.hashring.Add(fmt.Sprintf("%d", fragment_id)) +///////// SLICES ///////////////////////////////////////////////////////////////////////// - return &fragment +// A slice is the vertical combination of every fragment. +type Slice struct { + id int +} + +// Get a slice from a database +func (d *Database) GetSlice(slice_id int) (*Slice, error) { + d.mutex.Lock() + defer d.mutex.Unlock() + for _, slice := range d.slices { + if slice.id == slice_id { + return slice, nil + } + } + return nil, SliceDoesNotExistError +} + +// Add a slice to a database +func (d *Database) AddSlice(slice_id int) *Slice { + d.mutex.Lock() + defer d.mutex.Unlock() + slice := Slice{id: slice_id} + d.slices = append(d.slices, &slice) + // add intersections + for _, frame := range d.frames { + d.AddFrameSliceIntersect(frame, &slice) + } + return &slice +} + +func (d *Database) GetOrCreateSlice(slice_id int) *Slice { + d.mutex.Lock() + defer d.mutex.Unlock() + slice, err := d.GetSlice(slice_id) + if err == nil { + return slice + } + return d.AddSlice(slice_id) +} + + +///////// FRAME-SLICE INTERSECT ////////////////////////////////////////////////////////////// + +type FrameSliceIntersect struct { + frame *Frame + slice *Slice + fragments []*Fragment + hashring *consistent.Consistent } func (d *Database) AddFrameSliceIntersect(frame *Frame, slice *Slice) *FrameSliceIntersect { frameslice := FrameSliceIntersect{frame: frame, slice: slice} d.frame_slice_intersects = append(d.frame_slice_intersects, &frameslice) - frameslice.hashring = consistent.New() frameslice.hashring.NumberOfReplicas = 16 - return &frameslice } @@ -180,30 +234,125 @@ func (d *Database) GetFrameSliceIntersect(frame *Frame, slice *Slice) (*FrameSli return nil, FrameSliceIntersectDoesNotExistError } -// Get a frame from a database -func (d *Database) GetFrame(name string) (*Frame, error) { - for _, frame := range d.frames { - if frame.name == name { - return frame, nil +func (fsi *FrameSliceIntersect) GetFragment(fragment_id *uuid.UUID) (*Fragment, error) { + for _, fragment := range fsi.fragments { + if fragment.id == fragment_id { + return fragment, nil } } - return nil, FrameDoesNotExistError + return nil, FragmentDoesNotExistError } -// Get a slice from a database -func (d *Database) GetSlice(id int) (*Slice, error) { - for _, slice := range d.slices { - if slice.id == id { - return slice, nil - } - } - return nil, SliceDoesNotExistError +func (fsi *FrameSliceIntersect) AddFragment(fragment *Fragment) { + fsi.fragments = append(fsi.fragments, fragment) + fsi.hashring.Add(fragment.id.String()) } + + +///////// FRAGMENTS //////////////////////////////////////////////////////////////////////// + +// A fragment is a collection of bitmaps within a slice. The fragment contains a reference to the responsible node for that fragment. The node is in the form ip:port +type Fragment struct { + id *uuid.UUID + process *Process +} + +// rename this one +func (d *Database) OldGetFragment(bitmap Bitmap, profile_id int) (*Fragment, error) { + d.mutex.Lock() + defer d.mutex.Unlock() + slice, _ := d.GetSliceForProfile(profile_id) + frame, _ := d.GetFrame(bitmap.FrameType) + fsi, err := d.GetFrameSliceIntersect(frame, slice) + frag_id_s, err := fsi.hashring.Get(fmt.Sprintf("%d", bitmap.Id)) + frag_id, err := uuid.ParseHex(frag_id_s) + if err != nil { + log.Fatal(err) + } + return fsi.GetFragment(frag_id) +} + + +/* +// NOT IMPLEMENTED +// this would loop through all frame_slice_intersect[], then all fragmments to find a match +func (d *Database) GetFragment(fragment_id *uuid.UUID) *Fragment { +} +*/ +func (d *Database) GetFragment(frame *Frame, slice *Slice, fragment_id *uuid.UUID) (*Fragment, error) { + d.mutex.Lock() + defer d.mutex.Unlock() + fsi, err := d.GetFrameSliceIntersect(frame, slice) + if err != nil { + log.Fatal(err) + } + return fsi.GetFragment(fragment_id) +} + +func (d *Database) AddFragment(frame *Frame, slice *Slice, fragment_id *uuid.UUID) *Fragment { + d.mutex.Lock() + defer d.mutex.Unlock() + fsi, err := d.GetFrameSliceIntersect(frame, slice) + if err != nil { + log.Fatal(err) + } + fragment := Fragment{id: fragment_id} + fsi.AddFragment(&fragment) + return &fragment +} + +/* +func (d *Database) AllocateFragment(frame *Frame, slice *Slice) *Fragment { + // from ETCD, randomly get a process that has available_fragments > 0 + // atomically decrement available_fragments (as long as it's not 0) + // if it IS 0, try until we find a process with available capacity + + * + process, err := GetAvailableProcess() + if err != nil { + log.Fatal(err) + } + * + process_id, _ := uuid.NewV4() + process := NewProcess(process_id) + return nil + //return d.AddFragment(&frame, &slice, process) +} +*/ + +/* +func (d *Database) AddFragmentByProcess(frame *Frame, slice *Slice, process *Process) *Fragment { + d.mutex.Lock() + defer d.mutex.Unlock() + frameslice, _ := d.GetFrameSliceIntersect(frame, slice) + fragment_id, _ := uuid.NewV4() + fragment := Fragment{id: fragment_id, process: process} + frameslice.fragments = append(frameslice.fragments, &fragment) + frameslice.hashring.Add(fragment.id.String()) + return &fragment +} +*/ + +func (d *Database) GetOrCreateFragment(frame *Frame, slice *Slice, fragment_id *uuid.UUID) *Fragment { + d.mutex.Lock() + defer d.mutex.Unlock() + fragment, err := d.GetFragment(frame, slice, fragment_id) + if err == nil { + return fragment + } + return d.AddFragment(frame, slice, fragment_id) +} + +func (f *Fragment) SetProcess(process *Process) { + f.process = process +} + +/////////////////////////////////////////////////////////////////////////////////////////////// + + // Get a slice from a database func (d *Database) GetSliceForProfile(profile_id int) (*Slice, error) { - log.Println("GetSliceForProfile") - log.Println("profile_id:",profile_id) slice_id := profile_id / SLICE_WIDTH return d.GetSlice(slice_id) } @@ -213,18 +362,3 @@ type Bitmap struct { Id int FrameType string } - -func (d *Database) TestSetBit(bitmap Bitmap, profile_id int) { - log.Println("TestSetBit") - slice, _ := d.GetSliceForProfile(profile_id) - log.Println("slice:",slice) - frame, _ := d.GetFrame(bitmap.FrameType) - fsi, _ := d.GetFrameSliceIntersect(frame, slice) - frag_id_s,err := fsi.hashring.Get(fmt.Sprintf("%d", bitmap.Id)) - frag_id, err := strconv.Atoi(frag_id_s) - fragment, err := fsi.GetFragment(frag_id) - if err != nil { - log.Fatal(err) - } - spew.Dump(fragment) -} diff --git a/db/topology_test.go b/db/topology_test.go index aa7489739..f5ac55c00 100644 --- a/db/topology_test.go +++ b/db/topology_test.go @@ -3,6 +3,8 @@ package db import ( "testing" "log" + "github.com/nu7hatch/gouuid" + "github.com/davecgh/go-spew/spew" . "github.com/smartystreets/goconvey/convey" ) @@ -59,6 +61,16 @@ func TestTopology(t *testing.T) { database.AddFragment(frame, slice, loc1, 4) */ + uuid, _ := uuid.ParseHex("6a9aea17-2915-4eb4-858f-a8d7d4dc0a1e") + database.AddFragment(frame, slice, uuid) + /* + database.AddFragment(frame, slice, process) + database.AddFragment(frame, slice, process) + database.AddFragment(frame, slice, process) + database.AddFragment(frame, slice, process) + database.AddFragment(frame, slice, process) + */ + //fsi, _ := database.GetFrameSliceIntersect(frame, slice) //log.Println(fsi) @@ -74,10 +86,17 @@ func TestTopology(t *testing.T) { log.Println(errer) */ - //bitmap := Bitmap{Id: 555, FrameType: "general"} - //log.Println("bitmap:",bitmap) + bitmap := Bitmap{Id: 555, FrameType: "general"} + log.Println("bitmap:",bitmap) - //database.TestSetBit(bitmap, 65535) + /* + profile_id := 65535 + fragment, _ := database.OldGetFragment(bitmap, profile_id) + + spew.Dump("FRAGMENT") + spew.Dump(fragment) + */ + spew.Dump("DONE") }) }