From 3972d0cfaa1e4903d60b8ef758a83b8277bcbfd7 Mon Sep 17 00:00:00 2001 From: Michael Baird Date: Tue, 25 Apr 2017 13:31:12 -0500 Subject: [PATCH 01/10] added missing comments for godocs --- cache.go | 16 +++++++++++----- 1 file changed, 11 insertions(+), 5 deletions(-) diff --git a/cache.go b/cache.go index 1c338302e..a83dcdaf4 100644 --- a/cache.go +++ b/cache.go @@ -53,6 +53,7 @@ func NewLRUCache(maxEntries uint32) *LRUCache { return c } +// BulkAdd adds a count to the cache unsorted. You should Invalidate after completion. func (c *LRUCache) BulkAdd(id, n uint64) { c.Add(id, n) } @@ -194,6 +195,7 @@ func (c *RankCache) Invalidate() { c.invalidate() } +// Recalculate rebuilds the cache func (c *RankCache) Recalculate() { c.mu.Lock() defer c.mu.Unlock() @@ -303,19 +305,23 @@ type PairHeap struct { Pairs } +// Less implemets the Sort interface +// reports whether the element with index i should sort before the element with index j. func (p PairHeap) Less(i, j int) bool { return p.Pairs[i].Count < p.Pairs[j].Count } -func (h *Pairs) Push(x interface{}) { +// Push appends the element onto the Pair slice +func (p *Pairs) Push(x interface{}) { // Push and Pop use pointer receivers because they modify the slice's length, // not just its contents. - *h = append(*h, x.(Pair)) + *p = append(*p, x.(Pair)) } -func (h *Pairs) Pop() interface{} { - old := *h +// Pop removes the minimum element from the Pair slice +func (p *Pairs) Pop() interface{} { + old := *p n := len(old) x := old[n-1] - *h = old[0 : n-1] + *p = old[0 : n-1] return x } From 66496a73d718cb4b5b5e8209d44761e40c8534bb Mon Sep 17 00:00:00 2001 From: Michael Baird Date: Tue, 25 Apr 2017 13:57:20 -0500 Subject: [PATCH 02/10] additional comments for config --- cluster.go | 2 +- config.go | 15 +++++++++++---- 2 files changed, 12 insertions(+), 5 deletions(-) diff --git a/cluster.go b/cluster.go index 0ed3dfcb2..55c57125e 100644 --- a/cluster.go +++ b/cluster.go @@ -152,7 +152,7 @@ func (c *Cluster) NodeStates() map[string]string { return h } -// State returns the internal ClusterState representation. +// Status returns the internal ClusterState representation. func (c *Cluster) Status() *internal.ClusterStatus { return &internal.ClusterStatus{ Nodes: encodeClusterStatus(c.Nodes), diff --git a/config.go b/config.go index d5d725ece..a0cfb977f 100644 --- a/config.go +++ b/config.go @@ -3,10 +3,16 @@ package pilosa import "time" const ( - // DefaultHost is the default hostname and port to use. - DefaultHost = "localhost" - DefaultPort = "10101" - DefaultClusterType = "static" + // DefaultHost is the default hostname to use. + DefaultHost = "localhost" + + // DefaultPort is the default port use with the hostname. + DefaultPort = "10101" + + // DefaultClusterType sets the node intercommunication method + DefaultClusterType = "static" + + // DefaultInternalPort the port the nodes intercommunicate on DefaultInternalPort = "14000" ) @@ -67,6 +73,7 @@ func (d *Duration) UnmarshalText(text []byte) error { return nil } +// MarshalText writes duration value in text format. func (d Duration) MarshalText() (text []byte, err error) { return []byte(d.String()), nil } From 3512dd4549e2bdf001888249c9ad24d25258eb8a Mon Sep 17 00:00:00 2001 From: Michael Baird Date: Tue, 25 Apr 2017 14:15:02 -0500 Subject: [PATCH 03/10] missing comments --- handler.go | 2 ++ holder.go | 2 +- index.go | 2 ++ pilosa.go | 2 +- server.go | 4 ++-- stats.go | 4 ++-- 6 files changed, 10 insertions(+), 6 deletions(-) diff --git a/handler.go b/handler.go index 77dcb5700..b26fdefa8 100644 --- a/handler.go +++ b/handler.go @@ -58,6 +58,7 @@ func NewHandler() *Handler { return handler } +// NewRouter creates a Gorilla Mus http router func NewRouter(handler *Handler) *mux.Router { router := mux.NewRouter() router.HandleFunc("/index", handler.handleGetIndexes).Methods("GET") @@ -1315,6 +1316,7 @@ type QueryResponse struct { Err error } +// MarshalJSON marshals QueryResponse into a JSON-encoded byte slice func (resp *QueryResponse) MarshalJSON() ([]byte, error) { var output struct { Results []interface{} `json:"results,omitempty"` diff --git a/holder.go b/holder.go index 426817410..e8ef55483 100644 --- a/holder.go +++ b/holder.go @@ -356,7 +356,7 @@ type HolderSyncer struct { Closing <-chan struct{} } -// Returns true if the syncer has been marked to close. +// IsClosing returns true if the syncer has been marked to close. func (s *HolderSyncer) IsClosing() bool { select { case <-s.Closing: diff --git a/index.go b/index.go index 3fa8ac028..378f94ab7 100644 --- a/index.go +++ b/index.go @@ -250,6 +250,7 @@ func (i *Index) MaxSlice() uint64 { return max } +// SetRemoteMaxSlice sets the remote max slice value received from another node func (i *Index) SetRemoteMaxSlice(newmax uint64) { i.mu.Lock() defer i.mu.Unlock() @@ -273,6 +274,7 @@ func (i *Index) MaxInverseSlice() uint64 { return max } +// SetRemoteMaxInverseSlice sets the remote max inverse slice value received from another node func (i *Index) SetRemoteMaxInverseSlice(v uint64) { i.mu.Lock() defer i.mu.Unlock() diff --git a/pilosa.go b/pilosa.go index 16739fd7a..0b5532ff7 100644 --- a/pilosa.go +++ b/pilosa.go @@ -88,7 +88,7 @@ func decodeColumnAttrSet(pb *internal.ColumnAttrSet) *ColumnAttrSet { // TimeFormat is the go-style time format used to parse string dates. const TimeFormat = "2006-01-02T15:04" -// Restrict name using regex +// ValidateName restrict name using regex func ValidateName(name string) error { validName := nameRegexp.Match([]byte(name)) if validName == false { diff --git a/server.go b/server.go index 1b2f3fb87..e616d1bd9 100644 --- a/server.go +++ b/server.go @@ -281,13 +281,13 @@ func (s *Server) ReceiveMessage(pb proto.Message) error { return nil } -// Server implements StatusHandler. // LocalStatus returns the state of the local node as well as the // holder (indexes/frames) according to the local node. // In a gossip implementation, memberlist.Delegate.LocalState() uses this. +// Server implements StatusHandler. func (s *Server) LocalStatus() (proto.Message, error) { if s.Holder == nil { - return nil, errors.New("Server.Holder is nil.") + return nil, errors.New("Server.Holder is nil") } return &internal.NodeStatus{ Host: s.Host, diff --git a/stats.go b/stats.go index 1a6efe743..685a9ce3c 100644 --- a/stats.go +++ b/stats.go @@ -12,7 +12,7 @@ func init() { NopStatsClient = &nopStatsClient{} } -// Global expvar. +// Expvar Global expvar map. var Expvar = expvar.NewMap("index") // StatsClient represents a client to a stats server. @@ -39,9 +39,9 @@ type StatsClient interface { Timing(name string, value time.Duration) } +// NopStatsClient represents a client that doesn't do anything. var NopStatsClient StatsClient -// nopStatsClient represents a client that doesn't do anything. type nopStatsClient struct{} func (c *nopStatsClient) Tags() []string { return nil } From 35c1648c75a9c1d80ac6a4a9faa20bfad74363be Mon Sep 17 00:00:00 2001 From: Michael Baird Date: Tue, 25 Apr 2017 15:12:44 -0500 Subject: [PATCH 04/10] missing comments --- broadcast.go | 9 ++++++++- client.go | 29 +++++++++++++++-------------- executor.go | 4 ++-- 3 files changed, 25 insertions(+), 17 deletions(-) diff --git a/broadcast.go b/broadcast.go index 9e04d04f6..976527779 100644 --- a/broadcast.go +++ b/broadcast.go @@ -22,18 +22,22 @@ type StaticNodeSet struct { nodes []*Node } +// NewStaticNodeSet creates a static nodeset func NewStaticNodeSet() *StaticNodeSet { return &StaticNodeSet{} } +// Nodes implements the NodeSet interface and returns a list of nodes in the cluster func (s *StaticNodeSet) Nodes() []*Node { return s.nodes } +// Open implements the NodeSet interface to start network activity, but is a nop func (s *StaticNodeSet) Open() error { return nil } +// Join add nodes from the cluster to the NodeSet func (s *StaticNodeSet) Join(nodes []*Node) error { s.nodes = nodes return nil @@ -49,9 +53,9 @@ func init() { NopBroadcaster = &nopBroadcaster{} } +// NopBroadcaster represents a Broadcaster that doesn't do anything. var NopBroadcaster Broadcaster -// nopBroadcaster represents a Broadcaster that doesn't do anything. type nopBroadcaster struct{} // SendSync A no-op implemenetation of Broadcaster SendSync method. @@ -85,6 +89,7 @@ type nopBroadcastReceiver struct{} func (n *nopBroadcastReceiver) Start(b BroadcastHandler) error { return nil } +// NopBroadcastReceiver is a no op implementation of the BroadcastReceiver var NopBroadcastReceiver = &nopBroadcastReceiver{} const ( @@ -95,6 +100,7 @@ const ( MessageTypeDeleteFrame = 5 ) +// MarshalMessage encodes the protobuf message into a byte slice func MarshalMessage(m proto.Message) ([]byte, error) { var typ uint8 switch obj := m.(type) { @@ -118,6 +124,7 @@ func MarshalMessage(m proto.Message) ([]byte, error) { return append([]byte{typ}, buf...), nil } +// UnmarshalMessage decodes the byte slice into a protobuf message func UnmarshalMessage(buf []byte) (proto.Message, error) { typ, buf := buf[0], buf[1:] diff --git a/client.go b/client.go index 7606369ae..b4fc2c4d7 100644 --- a/client.go +++ b/client.go @@ -315,6 +315,7 @@ func (c *Client) Import(ctx context.Context, index, frame string, slice uint64, return nil } +// MarshalImportPayload marshalls the import parameters into a protobuf byte slice func MarshalImportPayload(index, frame string, slice uint64, bits []Bit) ([]byte, error) { // Separate row and column IDs to reduce allocations. rowIDs := Bits(bits).RowIDs() @@ -982,36 +983,36 @@ func (p Bits) Less(i, j int) bool { } // RowIDs returns a slice of all the row IDs. -func (a Bits) RowIDs() []uint64 { - other := make([]uint64, len(a)) - for i := range a { - other[i] = a[i].RowID +func (p Bits) RowIDs() []uint64 { + other := make([]uint64, len(p)) + for i := range p { + other[i] = p[i].RowID } return other } // ColumnIDs returns a slice of all the column IDs. -func (a Bits) ColumnIDs() []uint64 { - other := make([]uint64, len(a)) - for i := range a { - other[i] = a[i].ColumnID +func (p Bits) ColumnIDs() []uint64 { + other := make([]uint64, len(p)) + for i := range p { + other[i] = p[i].ColumnID } return other } // Timestamps returns a slice of all the timestamps. -func (a Bits) Timestamps() []int64 { - other := make([]int64, len(a)) - for i := range a { - other[i] = a[i].Timestamp +func (p Bits) Timestamps() []int64 { + other := make([]int64, len(p)) + for i := range p { + other[i] = p[i].Timestamp } return other } // GroupBySlice returns a map of bits by slice. -func (a Bits) GroupBySlice() map[uint64][]Bit { +func (p Bits) GroupBySlice() map[uint64][]Bit { m := make(map[uint64][]Bit) - for _, bit := range a { + for _, bit := range p { slice := bit.ColumnID / SliceWidth m[slice] = append(m[slice], bit) } diff --git a/executor.go b/executor.go index 3179770e3..26abd5808 100644 --- a/executor.go +++ b/executor.go @@ -801,7 +801,7 @@ func (e *Executor) executeSetRowAttrs(ctx context.Context, index string, c *pql. if err != nil { return fmt.Errorf("reading SetRowAttrs() row: %v", err) } else if !ok { - return fmt.Errorf("SetRowAttrs() row field '%v' required.", rowLabel) + return fmt.Errorf("SetRowAttrs() row field '%v' required", rowLabel) } // Copy args and remove reserved fields. @@ -860,7 +860,7 @@ func (e *Executor) executeBulkSetRowAttrs(ctx context.Context, index string, cal if err != nil { return nil, fmt.Errorf("reading SetRowAttrs() row: %v", rowLabel) } else if !ok { - return nil, fmt.Errorf("SetRowAttrs row field '%v' required.", rowLabel) + return nil, fmt.Errorf("SetRowAttrs row field '%v' required", rowLabel) } // Copy args and remove reserved fields. From fe0e8a3c5b5fc6552bbabb91510ca20411a6f94a Mon Sep 17 00:00:00 2001 From: Michael Baird Date: Tue, 25 Apr 2017 15:27:10 -0500 Subject: [PATCH 05/10] fixed enum comments --- broadcast.go | 1 + cluster.go | 4 +++- 2 files changed, 4 insertions(+), 1 deletion(-) diff --git a/broadcast.go b/broadcast.go index 976527779..b12dc0896 100644 --- a/broadcast.go +++ b/broadcast.go @@ -92,6 +92,7 @@ func (n *nopBroadcastReceiver) Start(b BroadcastHandler) error { return nil } // NopBroadcastReceiver is a no op implementation of the BroadcastReceiver var NopBroadcastReceiver = &nopBroadcastReceiver{} +// Broadcast message types const ( MessageTypeCreateSlice = 1 MessageTypeCreateIndex = 2 diff --git a/cluster.go b/cluster.go index 55c57125e..a6bcade81 100644 --- a/cluster.go +++ b/cluster.go @@ -13,8 +13,10 @@ const ( // DefaultReplicaN is the default number of replicas per partition. DefaultReplicaN = 1 +) - // NodeState represents node state returned in /status endpoint for a node in the cluster. +// NodeState represents node state returned in /status endpoint for a node in the cluster. +const ( NodeStateUp = "UP" NodeStateDown = "DOWN" ) From d2a78691423cdd95b8cf70cfbd68517069cb57ef Mon Sep 17 00:00:00 2001 From: Michael Baird Date: Tue, 25 Apr 2017 16:40:39 -0500 Subject: [PATCH 06/10] comment cleanup --- broadcast.go | 18 +++++++++--------- cache.go | 8 ++++---- client.go | 2 +- cluster.go | 4 ++-- config.go | 4 ++-- handler.go | 2 +- index.go | 4 ++-- pilosa.go | 2 +- stats.go | 2 +- 9 files changed, 23 insertions(+), 23 deletions(-) diff --git a/broadcast.go b/broadcast.go index b12dc0896..f36533a7a 100644 --- a/broadcast.go +++ b/broadcast.go @@ -17,27 +17,27 @@ type NodeSet interface { Open() error } -// StaticNodeSet represents a basic NodeSet for testing +// StaticNodeSet represents a basic NodeSet for testing. type StaticNodeSet struct { nodes []*Node } -// NewStaticNodeSet creates a static nodeset +// NewStaticNodeSet creates a statically defined NodeSet. func NewStaticNodeSet() *StaticNodeSet { return &StaticNodeSet{} } -// Nodes implements the NodeSet interface and returns a list of nodes in the cluster +// Nodes implements the NodeSet interface and returns a list of nodes in the cluster. func (s *StaticNodeSet) Nodes() []*Node { return s.nodes } -// Open implements the NodeSet interface to start network activity, but is a nop +// Open implements the NodeSet interface to start network activity, but for a static NodeSet it does nothing. func (s *StaticNodeSet) Open() error { return nil } -// Join add nodes from the cluster to the NodeSet +// Join sets the NodeSet nodes to the slices of Nodes passed in. func (s *StaticNodeSet) Join(nodes []*Node) error { s.nodes = nodes return nil @@ -89,10 +89,10 @@ type nopBroadcastReceiver struct{} func (n *nopBroadcastReceiver) Start(b BroadcastHandler) error { return nil } -// NopBroadcastReceiver is a no op implementation of the BroadcastReceiver +// NopBroadcastReceiver is a no-op implementation of the BroadcastReceiver. var NopBroadcastReceiver = &nopBroadcastReceiver{} -// Broadcast message types +// Broadcast message types. const ( MessageTypeCreateSlice = 1 MessageTypeCreateIndex = 2 @@ -101,7 +101,7 @@ const ( MessageTypeDeleteFrame = 5 ) -// MarshalMessage encodes the protobuf message into a byte slice +// MarshalMessage encodes the protobuf message into a byte slice. func MarshalMessage(m proto.Message) ([]byte, error) { var typ uint8 switch obj := m.(type) { @@ -125,7 +125,7 @@ func MarshalMessage(m proto.Message) ([]byte, error) { return append([]byte{typ}, buf...), nil } -// UnmarshalMessage decodes the byte slice into a protobuf message +// UnmarshalMessage decodes the byte slice into a protobuf message. func UnmarshalMessage(buf []byte) (proto.Message, error) { typ, buf := buf[0], buf[1:] diff --git a/cache.go b/cache.go index a83dcdaf4..cd172100b 100644 --- a/cache.go +++ b/cache.go @@ -195,7 +195,7 @@ func (c *RankCache) Invalidate() { c.invalidate() } -// Recalculate rebuilds the cache +// Recalculate rebuilds the cache. func (c *RankCache) Recalculate() { c.mu.Lock() defer c.mu.Unlock() @@ -305,18 +305,18 @@ type PairHeap struct { Pairs } -// Less implemets the Sort interface +// Less implemets the Sort interface. // reports whether the element with index i should sort before the element with index j. func (p PairHeap) Less(i, j int) bool { return p.Pairs[i].Count < p.Pairs[j].Count } -// Push appends the element onto the Pair slice +// Push appends the element onto the Pair slice. func (p *Pairs) Push(x interface{}) { // Push and Pop use pointer receivers because they modify the slice's length, // not just its contents. *p = append(*p, x.(Pair)) } -// Pop removes the minimum element from the Pair slice +// Pop removes the minimum element from the Pair slice. func (p *Pairs) Pop() interface{} { old := *p n := len(old) diff --git a/client.go b/client.go index b4fc2c4d7..2dfc6fadc 100644 --- a/client.go +++ b/client.go @@ -315,7 +315,7 @@ func (c *Client) Import(ctx context.Context, index, frame string, slice uint64, return nil } -// MarshalImportPayload marshalls the import parameters into a protobuf byte slice +// MarshalImportPayload marshalls the import parameters into a protobuf byte slice. func MarshalImportPayload(index, frame string, slice uint64, bits []Bit) ([]byte, error) { // Separate row and column IDs to reduce allocations. rowIDs := Bits(bits).RowIDs() diff --git a/cluster.go b/cluster.go index a6bcade81..e68e57926 100644 --- a/cluster.go +++ b/cluster.go @@ -127,7 +127,7 @@ func NewCluster() *Cluster { } } -// NodeSetHosts returns the list of host strings for NodeSet members +// NodeSetHosts returns the list of host strings for NodeSet members. func (c *Cluster) NodeSetHosts() []string { if c.NodeSet == nil { return []string{} @@ -154,7 +154,7 @@ func (c *Cluster) NodeStates() map[string]string { return h } -// Status returns the internal ClusterState representation. +// Status returns the internal ClusterStatus representation. func (c *Cluster) Status() *internal.ClusterStatus { return &internal.ClusterStatus{ Nodes: encodeClusterStatus(c.Nodes), diff --git a/config.go b/config.go index a0cfb977f..89fd3dc77 100644 --- a/config.go +++ b/config.go @@ -9,10 +9,10 @@ const ( // DefaultPort is the default port use with the hostname. DefaultPort = "10101" - // DefaultClusterType sets the node intercommunication method + // DefaultClusterType sets the node intercommunication method. DefaultClusterType = "static" - // DefaultInternalPort the port the nodes intercommunicate on + // DefaultInternalPort the port the nodes intercommunicate on. DefaultInternalPort = "14000" ) diff --git a/handler.go b/handler.go index b26fdefa8..57b05284b 100644 --- a/handler.go +++ b/handler.go @@ -58,7 +58,7 @@ func NewHandler() *Handler { return handler } -// NewRouter creates a Gorilla Mus http router +// NewRouter creates a Gorilla Mux http router. func NewRouter(handler *Handler) *mux.Router { router := mux.NewRouter() router.HandleFunc("/index", handler.handleGetIndexes).Methods("GET") diff --git a/index.go b/index.go index 378f94ab7..eb207a4a4 100644 --- a/index.go +++ b/index.go @@ -250,7 +250,7 @@ func (i *Index) MaxSlice() uint64 { return max } -// SetRemoteMaxSlice sets the remote max slice value received from another node +// SetRemoteMaxSlice sets the remote max slice value received from another node. func (i *Index) SetRemoteMaxSlice(newmax uint64) { i.mu.Lock() defer i.mu.Unlock() @@ -274,7 +274,7 @@ func (i *Index) MaxInverseSlice() uint64 { return max } -// SetRemoteMaxInverseSlice sets the remote max inverse slice value received from another node +// SetRemoteMaxInverseSlice sets the remote max inverse slice value received from another node. func (i *Index) SetRemoteMaxInverseSlice(v uint64) { i.mu.Lock() defer i.mu.Unlock() diff --git a/pilosa.go b/pilosa.go index 0b5532ff7..1736cccb9 100644 --- a/pilosa.go +++ b/pilosa.go @@ -88,7 +88,7 @@ func decodeColumnAttrSet(pb *internal.ColumnAttrSet) *ColumnAttrSet { // TimeFormat is the go-style time format used to parse string dates. const TimeFormat = "2006-01-02T15:04" -// ValidateName restrict name using regex +// ValidateName ensures that the name is a valid format. func ValidateName(name string) error { validName := nameRegexp.Match([]byte(name)) if validName == false { diff --git a/stats.go b/stats.go index 685a9ce3c..8cd599f2e 100644 --- a/stats.go +++ b/stats.go @@ -12,7 +12,7 @@ func init() { NopStatsClient = &nopStatsClient{} } -// Expvar Global expvar map. +// Expvar global expvar map. var Expvar = expvar.NewMap("index") // StatsClient represents a client to a stats server. From e2cab9b5a94e29b0ba87b48554ca2b22d0ad6200 Mon Sep 17 00:00:00 2001 From: Michael Baird Date: Tue, 25 Apr 2017 17:00:02 -0500 Subject: [PATCH 07/10] adding gossip comments --- gossip/gossip.go | 23 +++++++++++++++++------ 1 file changed, 17 insertions(+), 6 deletions(-) diff --git a/gossip/gossip.go b/gossip/gossip.go index c3ded776a..c3da20cfa 100644 --- a/gossip/gossip.go +++ b/gossip/gossip.go @@ -24,12 +24,13 @@ type GossipNodeSet struct { broadcasts *memberlist.TransmitLimitedQueue statusHandler pilosa.StatusHandler - config *GossipConfig + config *gossipConfig // The writer for any logging. LogOutput io.Writer } +// Nodes implements the NodeSet interface and returns a list of nodes in the cluster. func (g *GossipNodeSet) Nodes() []*pilosa.Node { a := make([]*pilosa.Node, 0, g.memberlist.NumMembers()) for _, n := range g.memberlist.Members() { @@ -38,11 +39,13 @@ func (g *GossipNodeSet) Nodes() []*pilosa.Node { return a } +// Start implements the BroadcastReceiver interface and starts listening for broadcast messages. func (g *GossipNodeSet) Start(h pilosa.BroadcastHandler) error { g.handler = h return nil } +// Open implements the NodeSet interface to start network activity. func (g *GossipNodeSet) Open() error { if g.handler == nil { return fmt.Errorf("opening GossipNodeSet: you must call Start(pilosa.BroadcastHandler) before calling Open()") @@ -75,7 +78,7 @@ func (g *GossipNodeSet) logger() *log.Logger { //////////////////////////////////////////////////////////////// -type GossipConfig struct { +type gossipConfig struct { gossipSeed string memberlistConfig *memberlist.Config } @@ -87,7 +90,7 @@ func NewGossipNodeSet(name string, gossipHost string, gossipPort int, gossipSeed } //TODO: pull memberlist config from pilosa.cfg file - g.config = &GossipConfig{ + g.config = &gossipConfig{ memberlistConfig: memberlist.DefaultLocalConfig(), gossipSeed: gossipSeed, } @@ -103,7 +106,7 @@ func NewGossipNodeSet(name string, gossipHost string, gossipPort int, gossipSeed return g } -// SendSync implementation of the Broadcaster interface +// SendSync implementation of the Broadcaster interface. func (g *GossipNodeSet) SendSync(pb proto.Message) error { msg, err := pilosa.MarshalMessage(pb) if err != nil { @@ -131,7 +134,7 @@ func (g *GossipNodeSet) SendSync(pb proto.Message) error { return eg.Wait() } -// SendAsync implementation of the Broadcaster interface +// SendAsync implementation of the Broadcaster interface. func (g *GossipNodeSet) SendAsync(pb proto.Message) error { msg, err := pilosa.MarshalMessage(pb) if err != nil { @@ -146,11 +149,13 @@ func (g *GossipNodeSet) SendAsync(pb proto.Message) error { return nil } -// implementation of the memberlist.Delegate interface +// NodeMeta implementation of the memberlist.Delegate interface. func (g *GossipNodeSet) NodeMeta(limit int) []byte { return []byte{} } +// NotifyMsg implementation of the memberlist.Delegate interface +// called when a user-data message is received. func (g *GossipNodeSet) NotifyMsg(b []byte) { m, err := pilosa.UnmarshalMessage(b) if err != nil { @@ -163,10 +168,14 @@ func (g *GossipNodeSet) NotifyMsg(b []byte) { } } +// GetBroadcasts implementation of the memberlist.Delegate interface +// called when user data messages can be broadcast. func (g *GossipNodeSet) GetBroadcasts(overhead, limit int) [][]byte { return g.broadcasts.GetBroadcasts(overhead, limit) } +// LocalState implementation of the memberlist.Delegate interface +// sends this Node's state data. func (g *GossipNodeSet) LocalState(join bool) []byte { pb, err := g.statusHandler.LocalStatus() if err != nil { @@ -183,6 +192,8 @@ func (g *GossipNodeSet) LocalState(join bool) []byte { return buf } +// MergeRemoteState implementation of the memberlist.Delegate interface +// receive and process the remote side side's LocalState. func (g *GossipNodeSet) MergeRemoteState(buf []byte, join bool) { // Unmarshal nodestate data. var pb internal.NodeStatus From 6ea72a1657752dc99b923339f3ff22fed7627df6 Mon Sep 17 00:00:00 2001 From: Michael Baird Date: Tue, 25 Apr 2017 17:20:26 -0500 Subject: [PATCH 08/10] added comments to http broadcaster --- httpbroadcast/messenger.go | 10 ++++++++-- 1 file changed, 8 insertions(+), 2 deletions(-) diff --git a/httpbroadcast/messenger.go b/httpbroadcast/messenger.go index ebd4b2e6a..566d39c22 100644 --- a/httpbroadcast/messenger.go +++ b/httpbroadcast/messenger.go @@ -62,11 +62,11 @@ func (h *HTTPBroadcaster) SendAsync(pb proto.Message) error { func (h *HTTPBroadcaster) nodes() ([]*pilosa.Node, error) { if h.server == nil { - return nil, errors.New("HTTPBroadcaster has no reference to Server.") + return nil, errors.New("HTTPBroadcaster has no reference to Server") } nodeset, ok := h.server.Cluster.NodeSet.(*HTTPNodeSet) if !ok { - return nil, errors.New("NodeSet cannot be caste to HTTPNodeSet.") + return nil, errors.New("NodeSet cannot be caste to HTTPNodeSet") } return nodeset.Nodes(), nil } @@ -106,12 +106,14 @@ func (h *HTTPBroadcaster) sendNodeMessage(node *pilosa.Node, msg []byte) error { return nil } +// HTTPBroadcastReceiver a Broadcaster that broadcasts messages over HTTP using the Pilosa Server. type HTTPBroadcastReceiver struct { port string handler pilosa.BroadcastHandler logOutput io.Writer } +// NewHTTPBroadcastReceiver returns a new instance of HTTPBroadcastReceiver. func NewHTTPBroadcastReceiver(port string, logOutput io.Writer) *HTTPBroadcastReceiver { return &HTTPBroadcastReceiver{ port: port, @@ -119,6 +121,7 @@ func NewHTTPBroadcastReceiver(port string, logOutput io.Writer) *HTTPBroadcastRe } } +// Start implements the BroadcastReceiver interface and starts listening for broadcast messages. func (rec *HTTPBroadcastReceiver) Start(b pilosa.BroadcastHandler) error { rec.handler = b go func() { @@ -166,14 +169,17 @@ func NewHTTPNodeSet() *HTTPNodeSet { return &HTTPNodeSet{} } +// Nodes implements the NodeSet interface and returns a list of nodes in the cluster. func (h *HTTPNodeSet) Nodes() []*pilosa.Node { return h.nodes } +// Open implements the NodeSet interface to start network activity, but for a HTTPNodeSet it does nothing. func (h *HTTPNodeSet) Open() error { return nil } +// Join sets the NodeSet nodes to the slices of Nodes passed in. func (h *HTTPNodeSet) Join(nodes []*pilosa.Node) error { h.nodes = nodes return nil From 2f876c46a67eed214d800d542181c342914e9112 Mon Sep 17 00:00:00 2001 From: Michael Baird Date: Tue, 25 Apr 2017 22:25:10 -0500 Subject: [PATCH 09/10] Broadcaster comment corrections --- broadcast.go | 2 +- gossip/gossip.go | 2 +- httpbroadcast/messenger.go | 4 ++-- 3 files changed, 4 insertions(+), 4 deletions(-) diff --git a/broadcast.go b/broadcast.go index f36533a7a..d4f75ebef 100644 --- a/broadcast.go +++ b/broadcast.go @@ -37,7 +37,7 @@ func (s *StaticNodeSet) Open() error { return nil } -// Join sets the NodeSet nodes to the slices of Nodes passed in. +// Join sets the NodeSet nodes to the slice of Nodes passed in. func (s *StaticNodeSet) Join(nodes []*Node) error { s.nodes = nodes return nil diff --git a/gossip/gossip.go b/gossip/gossip.go index c3da20cfa..a520548dc 100644 --- a/gossip/gossip.go +++ b/gossip/gossip.go @@ -39,7 +39,7 @@ func (g *GossipNodeSet) Nodes() []*pilosa.Node { return a } -// Start implements the BroadcastReceiver interface and starts listening for broadcast messages. +// Start implements the BroadcastReceiver interface and sets the BroadcastHandler func (g *GossipNodeSet) Start(h pilosa.BroadcastHandler) error { g.handler = h return nil diff --git a/httpbroadcast/messenger.go b/httpbroadcast/messenger.go index 566d39c22..92da723b9 100644 --- a/httpbroadcast/messenger.go +++ b/httpbroadcast/messenger.go @@ -106,7 +106,7 @@ func (h *HTTPBroadcaster) sendNodeMessage(node *pilosa.Node, msg []byte) error { return nil } -// HTTPBroadcastReceiver a Broadcaster that broadcasts messages over HTTP using the Pilosa Server. +// HTTPBroadcastReceiver handles incoming messages and pass them onto the implemented handler. type HTTPBroadcastReceiver struct { port string handler pilosa.BroadcastHandler @@ -179,7 +179,7 @@ func (h *HTTPNodeSet) Open() error { return nil } -// Join sets the NodeSet nodes to the slices of Nodes passed in. +// Join sets the NodeSet nodes to the slice of Nodes passed in. func (h *HTTPNodeSet) Join(nodes []*pilosa.Node) error { h.nodes = nodes return nil From ce399d722d78cc437599d3a82265af75c03ef271 Mon Sep 17 00:00:00 2001 From: Michael Baird Date: Wed, 26 Apr 2017 09:08:39 -0500 Subject: [PATCH 10/10] clarify HTTPBroadcastReceiver comment --- httpbroadcast/messenger.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/httpbroadcast/messenger.go b/httpbroadcast/messenger.go index 92da723b9..8537dfa45 100644 --- a/httpbroadcast/messenger.go +++ b/httpbroadcast/messenger.go @@ -106,7 +106,7 @@ func (h *HTTPBroadcaster) sendNodeMessage(node *pilosa.Node, msg []byte) error { return nil } -// HTTPBroadcastReceiver handles incoming messages and pass them onto the implemented handler. +// HTTPBroadcastReceiver unmarshals incoming messages over HTTP and passes them on to the handler. type HTTPBroadcastReceiver struct { port string handler pilosa.BroadcastHandler