diff --git a/api.go b/api.go index 66add8863..5b730b017 100644 --- a/api.go +++ b/api.go @@ -35,11 +35,10 @@ import ( // API provides the top level programmatic interface to Pilosa. It is usually // wrapped by a handler which provides an external interface (e.g. HTTP). type API struct { - Holder *Holder - Broadcaster Broadcaster - Cluster *Cluster - TranslateStore TranslateStore - server *Server + Holder *Holder + Broadcaster Broadcaster + Cluster *Cluster + server *Server } // APIOption is a functional option type for pilosa.API @@ -48,10 +47,9 @@ type APIOption func(*API) error func OptAPIServer(s *Server) APIOption { return func(a *API) error { a.server = s - a.TranslateStore = s.translateFile a.Holder = s.holder a.Broadcaster = s - a.Cluster = s.Cluster + a.Cluster = s.cluster return nil } } @@ -142,9 +140,9 @@ func (api *API) Query(ctx context.Context, req *QueryRequest) (QueryResponse, er } // Translate column attributes, if necessary. - if api.TranslateStore != nil { + if api.server.primaryTranslateStore != nil { for _, col := range resp.ColumnAttrSets { - v, err := api.TranslateStore.TranslateColumnToString(req.Index, col.ID) + v, err := api.server.primaryTranslateStore.TranslateColumnToString(req.Index, col.ID) if err != nil { return resp, err } @@ -710,7 +708,7 @@ func (api *API) LongQueryTime() time.Duration { if api.Cluster == nil { return 0 } - return api.Cluster.LongQueryTime + return api.Cluster.longQueryTime } func (api *API) indexField(indexName string, fieldName string, slice uint64) (*Index, *Field, error) { @@ -796,6 +794,48 @@ func (api *API) ResizeAbort() error { return errors.Wrap(err, "complete current job") } +// TranslateStoreBufferSize is the buffer size used for streaming data. +const TranslateStoreBufferSize = 65536 + +func (api *API) GetTranslateData(ctx context.Context, w io.WriteCloser, offset int64) error { + rc, err := api.server.primaryTranslateStore.Reader(ctx, offset) + if err != nil { + return errors.Wrap(err, "read from translate store") + } + + // Ensure reader is closed when the client disconnects. + go func() { <-ctx.Done(); rc.Close() }() + + go func() { + defer rc.Close() + defer w.Close() + + buf := make([]byte, TranslateStoreBufferSize) + + // Copy from reader to client until store or client disconnect. + for { + // Read from store. + n, err := rc.Read(buf) + if err == io.EOF { + return + } else if err != nil { + api.server.logger.Printf("api: translate store read error: %s", err) + return + } else if n == 0 { + continue + } + + // Write to response & flush. + if _, err := w.Write(buf[:n]); err != nil { + api.server.logger.Printf("api: translate store response write error: %s", err) + return + } + } + }() + + return nil +} + // State returns the cluster state which is usually "NORMAL", but could be // "STARTING", "RESIZING", or potentially others. See cluster.go for more // details. diff --git a/broadcast.go b/broadcast.go index 9b894fea2..b7b13fe04 100644 --- a/broadcast.go +++ b/broadcast.go @@ -23,30 +23,6 @@ import ( "github.com/pkg/errors" ) -// MemberSet represents an interface for Node membership and inter-node communication. -type MemberSet interface { - // Open starts any network activity implemented by the MemberSet - // Node is the local node, used for membership broadcasts. - Open(n *Node) error -} - -// StaticMemberSet represents a basic MemberSet for testing. -type StaticMemberSet struct { - nodes []*Node -} - -// NewStaticMemberSet creates a statically defined MemberSet. -func NewStaticMemberSet(nodes []*Node) *StaticMemberSet { - return &StaticMemberSet{ - nodes: nodes, - } -} - -// Open implements the MemberSet interface to start network activity, but for a static MemberSet it does nothing. -func (s *StaticMemberSet) Open(n *Node) error { - return nil -} - // Broadcaster is an interface for broadcasting messages. type Broadcaster interface { SendSync(pb proto.Message) error @@ -56,7 +32,6 @@ type Broadcaster interface { func init() { NopBroadcaster = &nopBroadcaster{} - NopGossiper = &nopGossiper{} } // NopBroadcaster represents a Broadcaster that doesn't do anything. @@ -85,39 +60,6 @@ type BroadcastHandler interface { ReceiveMessage(pb proto.Message) error } -// BroadcastReceiver is the interface for the object which will listen for and -// decode broadcast messages before passing them to pilosa to handle. The -// implementation of this could be an http server which listens for messages, -// gets the protobuf payload, and then passes it to -// BroadcastHandler.ReceiveMessage. -type BroadcastReceiver interface { - // Start starts listening for broadcast messages - it should return - // immediately, spawning a goroutine if necessary. - Start(BroadcastHandler) error -} - -type nopBroadcastReceiver struct{} - -func (n *nopBroadcastReceiver) Start(b BroadcastHandler) error { return nil } - -// NopBroadcastReceiver is a no-op implementation of the BroadcastReceiver. -var NopBroadcastReceiver = &nopBroadcastReceiver{} - -// Gossiper is an interface for sharing messages via gossip. -type Gossiper interface { - SendAsync(pb proto.Message) error -} - -// NopBroadcaster represents a Broadcaster that doesn't do anything. -var NopGossiper Gossiper - -type nopGossiper struct{} - -// SendAsync A no-op implementation of Gossiper SendAsync method. -func (n *nopGossiper) SendAsync(pb proto.Message) error { - return nil -} - // Broadcast message types. const ( messageTypeCreateSlice = iota diff --git a/cluster.go b/cluster.go index ae34bf03c..de728472a 100644 --- a/cluster.go +++ b/cluster.go @@ -212,28 +212,24 @@ type nodeAction struct { // Cluster represents a collection of nodes. type Cluster struct { - ID string - Node *Node - Nodes []*Node // TODO phase this out? - MemberSet MemberSet + id string + Node *Node + Nodes []*Node // TODO phase this out? // Hashing algorithm used to assign partitions to nodes. Hasher Hasher // The number of partitions in the cluster. - PartitionN int + partitionN int // The number of replicas a partition has. ReplicaN int // Threshold for logging long-running queries - LongQueryTime time.Duration + longQueryTime time.Duration // Maximum number of Set() or Clear() commands per request. - MaxWritesPerRequest int - - // EventReceiver receives NodeEvents pertaining to node membership. - EventReceiver EventReceiver + maxWritesPerRequest int // Data directory path. Path string @@ -243,8 +239,8 @@ type Cluster struct { Static bool // Static is primarily used for testing in a non-gossip environment. state string Coordinator string - Holder *Holder - Broadcaster Broadcaster + holder *Holder + broadcaster Broadcaster joiningLeavingNodes chan nodeAction @@ -261,7 +257,7 @@ type Cluster struct { wg sync.WaitGroup closing chan struct{} - Logger Logger + logger Logger InternalClient InternalClient } @@ -269,10 +265,9 @@ type Cluster struct { // NewCluster returns a new instance of Cluster with defaults. func NewCluster() *Cluster { return &Cluster{ - Hasher: &jmphasher{}, - PartitionN: DefaultPartitionN, - ReplicaN: 1, - EventReceiver: NopEventReceiver, + Hasher: &jmphasher{}, + partitionN: DefaultPartitionN, + ReplicaN: 1, joiningLeavingNodes: make(chan nodeAction, 10), // buffered channel jobs: make(map[int64]*resizeJob), @@ -281,7 +276,7 @@ func NewCluster() *Cluster { InternalClient: NewNopInternalClient(), - Logger: NopLogger, + logger: NopLogger, } } @@ -318,7 +313,7 @@ func (c *Cluster) setCoordinator(n *Node) error { _ = c.unprotectedUpdateCoordinator(n) c.mu.Unlock() // Send the update coordinator message to all nodes. - err := c.Broadcaster.SendSync( + err := c.broadcaster.SendSync( &internal.UpdateCoordinatorMessage{ New: EncodeNode(n), }) @@ -327,7 +322,7 @@ func (c *Cluster) setCoordinator(n *Node) error { } // Broadcast cluster status. - return c.Broadcaster.SendSync(c.Status()) + return c.broadcaster.SendSync(c.Status()) } // updateCoordinator updates this nodes Coordinator value as well as @@ -359,7 +354,7 @@ func (c *Cluster) unprotectedUpdateCoordinator(n *Node) bool { // addNode adds a node to the Cluster and updates and saves the // new topology. func (c *Cluster) addNode(node *Node) error { - c.Logger.Printf("add node %s to cluster on %s", node, c.Node) + c.logger.Printf("add node %s to cluster on %s", node, c.Node) // If the node being added is the coordinator, set it for this node. if node.IsCoordinator { @@ -410,13 +405,13 @@ func (c *Cluster) nodeIDs() []string { func (c *Cluster) setID(id string) { // Don't overwrite ClusterID. - if c.ID != "" { + if c.id != "" { return } - c.ID = id + c.id = id // Make sure the Topology is updated. - c.Topology.ClusterID = c.ID + c.Topology.ClusterID = c.id } func (c *Cluster) State() string { @@ -437,7 +432,7 @@ func (c *Cluster) setState(state string) { return } - c.Logger.Printf("change cluster state from %s to %s on %s", c.state, state, c.Node.ID) + c.logger.Printf("change cluster state from %s to %s on %s", c.state, state, c.Node.ID) var doCleanup bool @@ -457,13 +452,13 @@ func (c *Cluster) setState(state string) { if doCleanup { var cleaner HolderCleaner cleaner.Node = c.Node - cleaner.Holder = c.Holder + cleaner.Holder = c.holder cleaner.Cluster = c cleaner.Closing = c.closing // Clean holder. if err := cleaner.CleanHolder(); err != nil { - c.Logger.Printf("holder clean error: err=%s", err) + c.logger.Printf("holder clean error: err=%s", err) } } } @@ -479,7 +474,7 @@ func (c *Cluster) setNodeState(state string) error { State: state, } - c.Logger.Printf("Sending State %s (%s)", state, c.Coordinator) + c.logger.Printf("Sending State %s (%s)", state, c.Coordinator) if err := c.sendTo(c.coordinatorNode(), ns); err != nil { return fmt.Errorf("sending node state error: err=%s", err) } @@ -501,7 +496,7 @@ func (c *Cluster) receiveNodeState(nodeID string, state string) error { } c.Topology.nodeStates[nodeID] = state - c.Logger.Printf("received state %s (%s)", state, nodeID) + c.logger.Printf("received state %s (%s)", state, nodeID) // Set cluster state to NORMAL. if c.haveTopologyAgreement() && c.allNodesReady() { @@ -514,7 +509,7 @@ func (c *Cluster) receiveNodeState(nodeID string, state string) error { // Status returns the internal ClusterStatus representation. func (c *Cluster) Status() *internal.ClusterStatus { return &internal.ClusterStatus{ - ClusterID: c.ID, + ClusterID: c.id, State: c.state, Nodes: EncodeNodes(c.Nodes), } @@ -716,7 +711,7 @@ func (c *Cluster) fragSources(to *Cluster, idx *Index) (map[string][]*internal.R srcCluster = NewCluster() srcCluster.Nodes = Nodes(c.Nodes).Clone() srcCluster.Hasher = c.Hasher - srcCluster.PartitionN = c.PartitionN + srcCluster.partitionN = c.partitionN srcCluster.ReplicaN = 1 } @@ -786,7 +781,7 @@ func (c *Cluster) partition(index string, slice uint64) int { h := fnv.New64a() h.Write([]byte(index)) h.Write(buf[:]) - return int(h.Sum64() % uint64(c.PartitionN)) + return int(h.Sum64() % uint64(c.partitionN)) } // sliceNodes returns a list of nodes that own a fragment. @@ -861,7 +856,7 @@ func (h *jmphasher) Hash(key uint64, n int) int { return int(b) } -func (c *Cluster) open() error { +func (c *Cluster) setup() error { // Cluster always comes up in state STARTING until cluster membership is determined. c.state = ClusterStateStarting @@ -870,13 +865,13 @@ func (c *Cluster) open() error { return errors.Wrap(err, "loading topology") } - c.ID = c.Topology.ClusterID + c.id = c.Topology.ClusterID // Only the coordinator needs to consider the .topology file. if c.isCoordinator() { err := c.considerTopology() if err != nil { - return fmt.Errorf("considerTopology: %v", err) + return errors.Wrap(err, "considerTopology") } } @@ -885,12 +880,18 @@ func (c *Cluster) open() error { if err != nil { return errors.Wrap(err, "adding local node") } + return nil +} - // Open MemberSet communication. - if err := c.MemberSet.Open(c.Node); err != nil { - return fmt.Errorf("opening MemberSet: %v", err) +func (c *Cluster) open() error { + err := c.setup() + if err != nil { + return errors.Wrap(err, "setting up cluster") } + return c.waitForStarted() +} +func (c *Cluster) waitForStarted() error { // If not coordinator then wait for ClusterStatus from coordinator. if !c.isCoordinator() { // In the case where a node has been restarted and memberlist has @@ -905,13 +906,13 @@ func (c *Cluster) open() error { Event: uint32(NodeJoin), Node: EncodeNode(c.Node), } - if err := c.Broadcaster.SendSync(msg); err != nil { + if err := c.broadcaster.SendSync(msg); err != nil { return fmt.Errorf("sending restart NodeJoin: %v", err) } - c.Logger.Printf("%v wait for joining to complete", c.Node.ID) + c.logger.Printf("%v wait for joining to complete", c.Node.ID) <-c.joining - c.Logger.Printf("joining has completed") + c.logger.Printf("joining has completed") } return nil @@ -926,7 +927,7 @@ func (c *Cluster) close() error { } func (c *Cluster) markAsJoined() { - c.Logger.Printf("mark node as joined (received coordinator update)") + c.logger.Printf("mark node as joined (received coordinator update)") if !c.joined { c.joined = true close(c.joining) @@ -959,9 +960,9 @@ func (c *Cluster) allNodesReady() bool { func (c *Cluster) handleNodeAction(nodeAction nodeAction) error { j, err := c.generateResizeJob(nodeAction) if err != nil { - c.Logger.Printf("generateResizeJob error: err=%s", err) + c.logger.Printf("generateResizeJob error: err=%s", err) if err := c.setStateAndBroadcast(ClusterStateNormal); err != nil { - c.Logger.Printf("setStateAndBroadcast error: err=%s", err) + c.logger.Printf("setStateAndBroadcast error: err=%s", err) } return errors.Wrap(err, "setting state") } @@ -975,7 +976,7 @@ func (c *Cluster) handleNodeAction(nodeAction nodeAction) error { }) // Wait for the resizeJob to finish or be aborted. - c.Logger.Printf("wait for jobResult") + c.logger.Printf("wait for jobResult") jobResult := <-j.result // Make sure j.Run() didn't return an error. @@ -983,7 +984,7 @@ func (c *Cluster) handleNodeAction(nodeAction nodeAction) error { return errors.Wrap(err, "running job") } - c.Logger.Printf("received jobResult: %s", jobResult) + c.logger.Printf("received jobResult: %s", jobResult) switch jobResult { case resizeJobStateDone: if err := c.completeCurrentJob(resizeJobStateDone); err != nil { @@ -1009,12 +1010,12 @@ func (c *Cluster) setStateAndBroadcast(state string) error { return nil } // Broadcast cluster status changes to the cluster. - c.Logger.Printf("broadcasting ClusterStatus: %s", state) - return c.Broadcaster.SendSync(c.Status()) + c.logger.Printf("broadcasting ClusterStatus: %s", state) + return c.broadcaster.SendSync(c.Status()) } func (c *Cluster) sendTo(node *Node, msg proto.Message) error { - if err := c.Broadcaster.SendTo(node, msg); err != nil { + if err := c.broadcaster.SendTo(node, msg); err != nil { return errors.Wrap(err, "sending") } return nil @@ -1040,7 +1041,7 @@ func (c *Cluster) listenForJoins() { case nodeAction := <-c.joiningLeavingNodes: err := c.handleNodeAction(nodeAction) if err != nil { - c.Logger.Printf("handleNodeAction error: err=%s", err) + c.logger.Printf("handleNodeAction error: err=%s", err) continue } setNormal = true @@ -1052,7 +1053,7 @@ func (c *Cluster) listenForJoins() { if setNormal { // Put the cluster back to state NORMAL and broadcast. if err := c.setStateAndBroadcast(ClusterStateNormal); err != nil { - c.Logger.Printf("setStateAndBroadcast error: err=%s", err) + c.logger.Printf("setStateAndBroadcast error: err=%s", err) } } @@ -1063,7 +1064,7 @@ func (c *Cluster) listenForJoins() { case nodeAction := <-c.joiningLeavingNodes: err := c.handleNodeAction(nodeAction) if err != nil { - c.Logger.Printf("handleNodeAction error: err=%s", err) + c.logger.Printf("handleNodeAction error: err=%s", err) continue } setNormal = true @@ -1077,7 +1078,7 @@ func (c *Cluster) listenForJoins() { // added/removed. It also saves a reference to the resizeJob in the `jobs` map // for future lookup by JobID. func (c *Cluster) generateResizeJob(nodeAction nodeAction) (*resizeJob, error) { - c.Logger.Printf("generateResizeJob: %v", nodeAction) + c.logger.Printf("generateResizeJob: %v", nodeAction) c.mu.Lock() defer c.mu.Unlock() @@ -1085,7 +1086,7 @@ func (c *Cluster) generateResizeJob(nodeAction nodeAction) (*resizeJob, error) { if err != nil { return nil, errors.Wrap(err, "generating job") } - c.Logger.Printf("generated resizeJob: %d", j.ID) + c.logger.Printf("generated resizeJob: %d", j.ID) // Save job in jobs map for future reference. c.jobs[j.ID] = j @@ -1105,13 +1106,13 @@ func (c *Cluster) generateResizeJob(nodeAction nodeAction) (*resizeJob, error) { // the resize instructions to other nodes in the cluster. func (c *Cluster) generateResizeJobByAction(nodeAction nodeAction) (*resizeJob, error) { j := newResizeJob(c.Nodes, nodeAction.node, nodeAction.action) - j.Broadcaster = c.Broadcaster + j.Broadcaster = c.broadcaster // toCluster is a clone of Cluster with the new node added/removed for comparison. toCluster := NewCluster() toCluster.Nodes = Nodes(c.Nodes).Clone() toCluster.Hasher = c.Hasher - toCluster.PartitionN = c.PartitionN + toCluster.partitionN = c.partitionN toCluster.ReplicaN = c.ReplicaN if nodeAction.action == resizeJobActionRemove { toCluster.removeNodeBasicSorted(nodeAction.node) @@ -1127,7 +1128,7 @@ func (c *Cluster) generateResizeJobByAction(nodeAction nodeAction) (*resizeJob, } // Add to multiIndex the instructions for each index. - for _, idx := range c.Holder.Indexes() { + for _, idx := range c.holder.Indexes() { fragSources, err := c.fragSources(toCluster, idx) if err != nil { return nil, errors.Wrap(err, "getting sources") @@ -1149,7 +1150,7 @@ func (c *Cluster) generateResizeJobByAction(nodeAction nodeAction) (*resizeJob, Node: EncodeNode(toCluster.unprotectedNodeByID(id)), Coordinator: EncodeNode(c.coordinatorNode()), Sources: sources, - Schema: c.Holder.EncodeSchema(), // Include the schema to ensure it's in sync on the receiving node. + Schema: c.holder.EncodeSchema(), // Include the schema to ensure it's in sync on the receiving node. ClusterStatus: c.Status(), } j.Instructions = append(j.Instructions, instr) @@ -1176,21 +1177,21 @@ func (c *Cluster) completeCurrentJob(state string) error { // followResizeInstruction is run by any node that receives a ResizeInstruction. func (c *Cluster) followResizeInstruction(instr *internal.ResizeInstruction) error { - c.Logger.Printf("follow resize instruction on %s", c.Node.ID) + c.logger.Printf("follow resize instruction on %s", c.Node.ID) // Make sure the cluster status on this node agrees with the Coordinator // before attempting a resize. if err := c.mergeClusterStatus(instr.ClusterStatus); err != nil { return errors.Wrap(err, "merging cluster status") } - c.Logger.Printf("MergeClusterStatus done, start goroutine") + c.logger.Printf("MergeClusterStatus done, start goroutine") // The actual resizing runs in a goroutine because we don't want to block // the distribution of other ResizeInstructions to the rest of the cluster. go func() { // Make sure the holder has opened. - <-c.Holder.opened + <-c.holder.opened // Prepare the return message. complete := &internal.ResizeInstructionComplete{ @@ -1203,19 +1204,19 @@ func (c *Cluster) followResizeInstruction(instr *internal.ResizeInstruction) err if err := func() error { // Sync the schema received in the resize instruction. - c.Logger.Printf("Holder ApplySchema") - if err := c.Holder.ApplySchema(instr.Schema); err != nil { + c.logger.Printf("Holder ApplySchema") + if err := c.holder.ApplySchema(instr.Schema); err != nil { return errors.Wrap(err, "applying schema") } // Request each source file in ResizeSources. for _, src := range instr.Sources { - c.Logger.Printf("get slice %d for index %s from host %s", src.Slice, src.Index, src.Node.URI) + c.logger.Printf("get slice %d for index %s from host %s", src.Slice, src.Index, src.Node.URI) srcURI := decodeURI(src.Node.URI) // Retrieve field. - f := c.Holder.Field(src.Index, src.Field) + f := c.holder.Field(src.Index, src.Field) if f == nil { return ErrFieldNotFound } @@ -1233,7 +1234,7 @@ func (c *Cluster) followResizeInstruction(instr *internal.ResizeInstruction) err } // Stream slice from remote node. - c.Logger.Printf("retrieve slice %d for index %s from host %s", src.Slice, src.Index, src.Node.URI) + c.logger.Printf("retrieve slice %d for index %s from host %s", src.Slice, src.Index, src.Node.URI) rd, err := c.InternalClient.RetrieveSliceFromURI(context.Background(), src.Index, src.Field, src.Slice, srcURI) if err != nil { // For now it is an acceptable error if the fragment is not found @@ -1265,7 +1266,7 @@ func (c *Cluster) followResizeInstruction(instr *internal.ResizeInstruction) err } if err := c.sendTo(DecodeNode(instr.Coordinator), complete); err != nil { - c.Logger.Printf("sending resizeInstructionComplete error: err=%s", err) + c.logger.Printf("sending resizeInstructionComplete error: err=%s", err) } }() return nil @@ -1580,10 +1581,10 @@ func decodeTopology(topology *internal.Topology) (*Topology, error) { func (c *Cluster) considerTopology() error { // Create ClusterID if one does not already exist. - if c.ID == "" { + if c.id == "" { u := uuid.NewV4() - c.ID = u.String() - c.Topology.ClusterID = c.ID + c.id = u.String() + c.Topology.ClusterID = c.id } if c.Static { @@ -1619,7 +1620,7 @@ func (c *Cluster) ReceiveEvent(e *NodeEvent) error { switch e.Event { case NodeJoin: - c.Logger.Printf("received NodeJoin event: %v", e) + c.logger.Printf("received NodeJoin event: %v", e) // Ignore the event if this is not the coordinator. if !c.isCoordinator() { return nil @@ -1639,7 +1640,7 @@ func (c *Cluster) nodeJoin(node *Node) error { // A host that is not part of the topology can't be added to the STARTING cluster. if !c.Topology.ContainsID(node.ID) { err := fmt.Sprintf("host is not in topology: %s", node.ID) - c.Logger.Printf("%v", err) + c.logger.Printf("%v", err) return errors.New(err) } @@ -1650,7 +1651,7 @@ func (c *Cluster) nodeJoin(node *Node) error { // Only change to normal if there is no existing data. Otherwise, // the coordinator needs to wait to receive READY messages (nodeStates) // from remote nodes before setting the cluster to state NORMAL. - if ok, err := c.Holder.HasData(); !ok && err == nil { + if ok, err := c.holder.HasData(); !ok && err == nil { // If the result of the previous AddNode completed the joining of nodes // in the topology, then change the state to NORMAL. if c.haveTopologyAgreement() { @@ -1678,7 +1679,7 @@ func (c *Cluster) nodeJoin(node *Node) error { } // If the holder does not yet contain data, go ahead and add the node. - if ok, err := c.Holder.HasData(); !ok && err == nil { + if ok, err := c.holder.HasData(); !ok && err == nil { if err := c.addNode(node); err != nil { return errors.Wrap(err, "adding node") } @@ -1732,7 +1733,7 @@ func (c *Cluster) nodeLeave(node *Node) error { } // If the holder does not yet contain data, go ahead and remove the node. - if ok, err := c.Holder.HasData(); !ok && err == nil { + if ok, err := c.holder.HasData(); !ok && err == nil { if err := c.removeNode(n); err != nil { return errors.Wrap(err, "removing node") } @@ -1754,7 +1755,7 @@ func (c *Cluster) nodeLeave(node *Node) error { func (c *Cluster) mergeClusterStatus(cs *internal.ClusterStatus) error { c.mu.Lock() defer c.mu.Unlock() - c.Logger.Printf("merge cluster status: %v", cs) + c.logger.Printf("merge cluster status: %v", cs) // Ignore status updates from self (coordinator). if c.unprotectedIsCoordinator() { return nil @@ -1810,6 +1811,5 @@ func (c *Cluster) setStatic(hosts []string) error { } c.Nodes = append(c.Nodes, &Node{URI: *uri}) } - c.MemberSet = NewStaticMemberSet(c.Nodes) return nil } diff --git a/cluster_internal_test.go b/cluster_internal_test.go index bdd51c047..b53180b39 100644 --- a/cluster_internal_test.go +++ b/cluster_internal_test.go @@ -25,6 +25,7 @@ import ( "github.com/davecgh/go-spew/spew" "github.com/pilosa/pilosa/internal" + "github.com/pkg/errors" ) // Ensure that fragCombos creates the correct fragment mapping. @@ -340,7 +341,7 @@ func TestCluster_Owners(t *testing.T) { func TestCluster_Partition(t *testing.T) { if err := quick.Check(func(index string, slice uint64, partitionN int) bool { c := NewCluster() - c.PartitionN = partitionN + c.partitionN = partitionN partitionID := c.partition(index, slice) if partitionID < 0 || partitionID >= partitionN { @@ -591,10 +592,10 @@ func TestCluster_ResizeStates(t *testing.T) { tc.WriteTopology(node.Path, top) // Open TestCluster. - expected := "considerTopology: coordinator node0 is not in topology: [some-other-host]" + expected := "coordinator node0 is not in topology: [some-other-host]" err := tc.Open() - if err == nil || err.Error() != expected { - t.Errorf("did not receive expected error: %s", expected) + if err == nil || errors.Cause(err).Error() != expected { + t.Errorf("did not receive expected error, got: %s", errors.Cause(err).Error()) } // Close TestCluster. @@ -704,7 +705,7 @@ func TestCluster_ResizeStates(t *testing.T) { // Before starting the resize, get the CheckSum to use for // comparison later. - node0Field := node0.Holder.Field("i", "f") + node0Field := node0.holder.Field("i", "f") node0View := node0Field.View("standard") node0Fragment := node0View.Fragment(1) node0Checksum := node0Fragment.Checksum() @@ -733,7 +734,7 @@ func TestCluster_ResizeStates(t *testing.T) { // Bits // Verify that node-1 contains the fragment (i/f/standard/1) transferred from node-0. - node1Field := node1.Holder.Field("i", "f") + node1Field := node1.holder.Field("i", "f") node1View := node1Field.View("standard") node1Fragment := node1View.Fragment(1) diff --git a/ctl/import_test.go b/ctl/import_test.go index 5500fdadf..522de0eba 100644 --- a/ctl/import_test.go +++ b/ctl/import_test.go @@ -188,7 +188,7 @@ func TestImportCommand_BugOverwriteValue(t *testing.T) { t.Fatal(err) } - cm.Host = cmd.Server.Addr().String() + cm.Host = cmd.Server.URI.HostPort() http.DefaultClient.Do(MustNewHTTPRequest("POST", "http://"+cm.Host+"/index/i", strings.NewReader(""))) http.DefaultClient.Do(MustNewHTTPRequest("POST", "http://"+cm.Host+"/index/i/field/f", strings.NewReader(`{"options":{"type": "int", "min": 0, "max":2147483648 }}`))) diff --git a/event.go b/event.go index 5df69361b..69229b1e4 100644 --- a/event.go +++ b/event.go @@ -35,21 +35,3 @@ type NodeEvent struct { type EventHandler interface { ReceiveEvent(e *NodeEvent) error } - -// EventReceiver is the interface for the object which will listen for and -// decode broadcast messages before passing them to pilosa to handle. The -// implementation of this could be an http server which listens for messages, -// gets the protobuf payload, and then passes it to -// EventHandler.ReceiveMessage. -type EventReceiver interface { - // Start starts listening for broadcast messages - it should return - // immediately, spawning a goroutine if necessary. - Start(EventHandler) error -} - -type nopEventReceiver struct{} - -func (n *nopEventReceiver) Start(e EventHandler) error { return nil } - -// NopEventReceiver is a no-op implementation of the EventReceiver. -var NopEventReceiver = &nopEventReceiver{} diff --git a/fragment.go b/fragment.go index 4fc24d71e..9b4780795 100644 --- a/fragment.go +++ b/fragment.go @@ -1865,7 +1865,7 @@ func (s *FragmentSyncer) syncBlock(id int) error { // Generate query with sets & clears, and group the requests to not exceed MaxWritesPerRequest. total := len(set.columnIDs) + len(clear.columnIDs) - maxWrites := s.Cluster.MaxWritesPerRequest + maxWrites := s.Cluster.maxWritesPerRequest if maxWrites <= 0 { maxWrites = 5000 } diff --git a/gossip/gossip.go b/gossip/gossip.go index 4749c9c58..0355260d6 100644 --- a/gossip/gossip.go +++ b/gossip/gossip.go @@ -38,7 +38,6 @@ var _ memberlist.Delegate = &GossipMemberSet{} // GossipMemberSet represents a gossip implementation of MemberSet using memberlist. type GossipMemberSet struct { mu sync.RWMutex - node *pilosa.Node memberlist *memberlist.Memberlist handler pilosa.BroadcastHandler @@ -63,7 +62,7 @@ func (g *GossipMemberSet) GetBindAddr() string { } // Open implements the MemberSet interface to start network activity. -func (g *GossipMemberSet) Open(n *pilosa.Node) error { +func (g *GossipMemberSet) Open() error { err := g.gossipEventReceiver.Start(g.pserver) if err != nil { return errors.Wrap(err, "starting event delegate") @@ -72,8 +71,6 @@ func (g *GossipMemberSet) Open(n *pilosa.Node) error { return fmt.Errorf("must call Start(pilosa.BroadcastHandler) before calling Open()") } - g.node = n - g.mu.Lock() g.memberlist, err = memberlist.Create(g.config.memberlistConfig) g.mu.Unlock() @@ -164,7 +161,8 @@ func WithLogger(logger *log.Logger) GossipMemberSetOption { } // NewGossipMemberSet returns a new instance of GossipMemberSet based on options. -func NewGossipMemberSet(name string, host string, cfg Config, s *pilosa.Server, options ...GossipMemberSetOption) (*GossipMemberSet, error) { +func NewGossipMemberSet(cfg Config, s *pilosa.Server, options ...GossipMemberSetOption) (*GossipMemberSet, error) { + host := s.Node().URI.Host() g := &GossipMemberSet{ Logger: pilosa.NopLogger, } @@ -209,11 +207,11 @@ func NewGossipMemberSet(name string, host string, cfg Config, s *pilosa.Server, // memberlist config conf := memberlist.DefaultWANConfig() conf.Transport = g.transport.Net - conf.Name = name - conf.BindAddr = host + conf.Name = s.Node().ID + conf.BindAddr = s.Node().URI.Host() conf.BindPort = port conf.AdvertisePort = port - conf.AdvertiseAddr = hostToIP(host) + conf.AdvertiseAddr = hostToIP(s.Node().URI.Host()) // conf.TCPTimeout = time.Duration(cfg.StreamTimeout) conf.SuspicionMult = cfg.SuspicionMult @@ -241,7 +239,7 @@ func NewGossipMemberSet(name string, host string, cfg Config, s *pilosa.Server, // NodeMeta implementation of the memberlist.Delegate interface. func (g *GossipMemberSet) NodeMeta(limit int) []byte { - buf, err := proto.Marshal(pilosa.EncodeNode(g.node)) + buf, err := proto.Marshal(pilosa.EncodeNode(g.pserver.Node())) if err != nil { g.Logger.Printf("marshal message error: %s", err) return []byte{} diff --git a/http/client_test.go b/http/client_test.go index 5ac29ec2c..3ca4cde50 100644 --- a/http/client_test.go +++ b/http/client_test.go @@ -220,7 +220,7 @@ func TestClient_MultiNode(t *testing.T) { // Ensure client can bulk import data. func TestClient_Import(t *testing.T) { cmd := test.MustRunMainWithCluster(t, 1)[0] - host := cmd.Server.Addr().String() + host := cmd.URL() holder := cmd.Server.Holder() hldr := test.Holder{Holder: holder} @@ -250,7 +250,7 @@ func TestClient_Import(t *testing.T) { // Ensure client can bulk import value data. func TestClient_ImportValue(t *testing.T) { cmd := test.MustRunMainWithCluster(t, 1)[0] - host := cmd.Server.Addr().String() + host := cmd.URL() holder := cmd.Server.Holder() hldr := test.Holder{Holder: holder} @@ -330,7 +330,7 @@ func TestClient_FragmentBlocks(t *testing.T) { // Set a bit on a different slice. hldr.SetBit("i", "f", 0, 1) - c := MustNewClient(cmd.Server.Addr().String(), defaultClient) + c := MustNewClient(cmd.URL(), defaultClient) blocks, err := c.FragmentBlocks(context.Background(), nil, "i", "f", 0) if err != nil { t.Fatal(err) diff --git a/http/handler.go b/http/handler.go index bbe784f69..58a3f751d 100644 --- a/http/handler.go +++ b/http/handler.go @@ -1355,25 +1355,22 @@ func (h *Handler) GetAPI() *pilosa.API { type defaultClusterMessageResponse struct{} -// TranslateStoreBufferSize is the buffer size used for streaming data. -const TranslateStoreBufferSize = 65536 - func (h *Handler) handleGetTranslateData(w http.ResponseWriter, r *http.Request) { q := r.URL.Query() offset, _ := strconv.ParseInt(q.Get("offset"), 10, 64) - rc, err := h.API.TranslateStore.Reader(r.Context(), offset) - if err == pilosa.ErrNotImplemented { - http.Error(w, err.Error(), http.StatusNotImplemented) - return - } else if err != nil { - http.Error(w, err.Error(), http.StatusInternalServerError) + pipeR, pipeW := io.Pipe() + + err := h.API.GetTranslateData(r.Context(), pipeW, offset) + + if err != nil { + if errors.Cause(err) == pilosa.ErrNotImplemented { + http.Error(w, err.Error(), http.StatusNotImplemented) + } else { + http.Error(w, err.Error(), http.StatusInternalServerError) + } return } - defer rc.Close() - - // Ensure reader is closed when the client disconnects. - go func() { <-r.Context().Done(); rc.Close() }() // Flush header so client can continue. w.WriteHeader(http.StatusOK) @@ -1381,28 +1378,7 @@ func (h *Handler) handleGetTranslateData(w http.ResponseWriter, r *http.Request) w.Flush() } - // Copy from reader to client until store or client disconnect. - buf := make([]byte, TranslateStoreBufferSize) - for { - // Read from store. - n, err := rc.Read(buf) - if err == io.EOF { - return - } else if err != nil { - h.Logger.Printf("http: translate store read error: %s", err) - return - } else if n == 0 { - continue - } - - // Write to response & flush. - if _, err := w.Write(buf[:n]); err != nil { - h.Logger.Printf("http: translate store response write error: %s", err) - return - } else if w, ok := w.(http.Flusher); ok { - w.Flush() - } - } + io.Copy(w, pipeR) } type queryValidationSpec struct { diff --git a/http/translator_test.go b/http/translator_test.go index 15aef62a4..d317944e4 100644 --- a/http/translator_test.go +++ b/http/translator_test.go @@ -4,6 +4,7 @@ import ( "context" "io" "io/ioutil" + gohttp "net/http" "testing" "time" @@ -15,8 +16,6 @@ import ( ) func TestTranslateStore_Reader(t *testing.T) { - t.Skip() // Until test.NewServer() works - // Ensure client can connect and stream the translate store data. t.Run("OK", func(t *testing.T) { t.Run("ServerDisconnect", func(t *testing.T) { @@ -46,15 +45,30 @@ func TestTranslateStore_Reader(t *testing.T) { // Setup handler on test server. var translateStore mock.TranslateStore + translateStore.ReaderFunc = func(ctx context.Context, off int64) (io.ReadCloser, error) { - if off != 100 { - t.Fatalf("unexpected off: %d", off) + // Check context to make sure this is the call we are looking for. + // (Something else calls ReaderFunc on server startup) + if ctx.Value(gohttp.ServerContextKey) != nil { + if off != 100 { + t.Fatalf("unexpected off: %d", off) + } + return &mrc, nil } - return &mrc, nil + mrc2 := mock.ReadCloser{ + ReadFunc: func(p []byte) (int, error) { + return 0, io.EOF + }, + CloseFunc: func() error { + return nil + }, + } + return &mrc2, nil } opts := server.OptCommandServerOptions(pilosa.OptServerPrimaryTranslateStore(translateStore)) main := test.MustRunMainWithCluster(t, 1, []server.CommandOption{opts})[0] + defer main.Close() // Connect to server and stream all available data. @@ -128,6 +142,7 @@ func TestTranslateStore_Reader(t *testing.T) { opts := server.OptCommandServerOptions(pilosa.OptServerPrimaryTranslateStore(translateStore)) main := test.MustRunMainWithCluster(t, 1, []server.CommandOption{opts})[0] + defer main.Close() _, err := http.NewTranslateStore(main.Server.URI.String()).Reader(context.Background(), 0) if err != pilosa.ErrNotImplemented { diff --git a/server.go b/server.go index 74bedab5b..c07720a54 100644 --- a/server.go +++ b/server.go @@ -18,7 +18,6 @@ import ( "context" "fmt" "log" - "net" "os" "os/exec" "path/filepath" @@ -53,7 +52,7 @@ type Server struct { // Internal holder *Holder - Cluster *Cluster + cluster *Cluster translateFile *TranslateFile diagnostics *DiagnosticsCollector executor *Executor @@ -65,12 +64,13 @@ type Server struct { gcNotifier GCNotifier logger Logger - NodeID string + nodeID string URI URI antiEntropyInterval time.Duration metricInterval time.Duration diagnosticInterval time.Duration maxWritesPerRequest int + isCoordinator bool primaryTranslateStore TranslateStore @@ -95,7 +95,7 @@ func OptServerLogger(l Logger) ServerOption { func OptServerReplicaN(n int) ServerOption { return func(s *Server) error { - s.Cluster.ReplicaN = n + s.cluster.ReplicaN = n return nil } } @@ -123,7 +123,7 @@ func OptServerAntiEntropyInterval(interval time.Duration) ServerOption { func OptServerLongQueryTime(dur time.Duration) ServerOption { return func(s *Server) error { - s.Cluster.LongQueryTime = dur + s.cluster.longQueryTime = dur return nil } } @@ -160,7 +160,7 @@ func OptServerInternalClient(c InternalClient) ServerOption { return func(s *Server) error { s.executor = NewExecutor(OptExecutorInternalQueryClient(c)) s.defaultClient = c - s.Cluster.InternalClient = c + s.cluster.InternalClient = c return nil } } @@ -203,11 +203,18 @@ func OptServerClusterDisabled(disabled bool, hosts []string) ServerOption { } } +func OptServerIsCoordinator(is bool) ServerOption { + return func(s *Server) error { + s.isCoordinator = is + return nil + } +} + // NewServer returns a new instance of Server. func NewServer(opts ...ServerOption) (*Server, error) { s := &Server{ closing: make(chan struct{}), - Cluster: NewCluster(), + cluster: NewCluster(), holder: NewHolder(), diagnostics: NewDiagnosticsCollector(DefaultDiagnosticServer), systemInfo: NewNopSystemInfo(), @@ -238,9 +245,9 @@ func NewServer(opts ...ServerOption) (*Server, error) { s.holder.Logger = s.logger s.holder.Stats.SetLogger(s.logger) - s.Cluster.Path = path - s.Cluster.Logger = s.logger - s.Cluster.Holder = s.holder + s.cluster.Path = path + s.cluster.logger = s.logger + s.cluster.holder = s.holder // Initialize translation database. s.translateFile = NewTranslateFile() @@ -248,29 +255,41 @@ func NewServer(opts ...ServerOption) (*Server, error) { s.translateFile.PrimaryTranslateStore = s.primaryTranslateStore // Get or create NodeID. - s.NodeID = s.LoadNodeID() + s.nodeID = s.loadNodeID() + if s.isCoordinator { + s.cluster.Coordinator = s.nodeID + } + // Set Cluster Node. node := &Node{ - ID: s.NodeID, + ID: s.nodeID, URI: s.URI, - IsCoordinator: s.Cluster.Coordinator == s.NodeID, + IsCoordinator: s.cluster.Coordinator == s.nodeID, } - s.Cluster.Node = node + s.cluster.Node = node if s.clusterDisabled { - err := s.Cluster.setStatic(s.hosts) + err := s.cluster.setStatic(s.hosts) if err != nil { return nil, errors.Wrap(err, "setting cluster static") } } // Append the NodeID tag to stats. - s.holder.Stats = s.holder.Stats.WithTags(fmt.Sprintf("NodeID:%s", s.NodeID)) + s.holder.Stats = s.holder.Stats.WithTags(fmt.Sprintf("NodeID:%s", s.nodeID)) s.executor.Holder = s.holder s.executor.Node = node - s.executor.Cluster = s.Cluster + s.executor.Cluster = s.cluster s.executor.TranslateStore = s.translateFile s.executor.MaxWritesPerRequest = s.maxWritesPerRequest + s.cluster.broadcaster = s + s.cluster.maxWritesPerRequest = s.maxWritesPerRequest + s.holder.Broadcaster = s + + err = s.cluster.setup() + if err != nil { + return nil, errors.Wrap(err, "setting up cluster") + } return s, nil } @@ -290,15 +309,8 @@ func (s *Server) Open() error { return err } - // Cluster settings. - s.Cluster.Broadcaster = s - s.Cluster.MaxWritesPerRequest = s.maxWritesPerRequest - - // Initialize Holder. - s.holder.Broadcaster = s - // Open Cluster management. - if err := s.Cluster.open(); err != nil { + if err := s.cluster.waitForStarted(); err != nil { return fmt.Errorf("opening Cluster: %v", err) } @@ -306,7 +318,7 @@ func (s *Server) Open() error { if err := s.holder.Open(); err != nil { return fmt.Errorf("opening Holder: %v", err) } - if err := s.Cluster.setNodeState(NodeStateReady); err != nil { + if err := s.cluster.setNodeState(NodeStateReady); err != nil { return fmt.Errorf("setting nodeState: %v", err) } @@ -315,7 +327,7 @@ func (s *Server) Open() error { // the cluster without waiting for data to load on the coordinator. Before // this starts, the joins are queued up in the Cluster.joiningLeavingNodes // buffered channel. - s.Cluster.listenForJoins() + s.cluster.listenForJoins() // Start background monitoring. s.wg.Add(3) @@ -332,8 +344,8 @@ func (s *Server) Close() error { close(s.closing) s.wg.Wait() - if s.Cluster != nil { - s.Cluster.close() + if s.cluster != nil { + s.cluster.close() } if s.holder != nil { s.holder.Close() @@ -345,37 +357,20 @@ func (s *Server) Close() error { return nil } -// LoadNodeID gets NodeID from disk, or creates a new value. +// loadNodeID gets NodeID from disk, or creates a new value. // If server.NodeID is already set, a new ID is not created. -func (s *Server) LoadNodeID() string { - if s.NodeID != "" { - return s.NodeID +func (s *Server) loadNodeID() string { + if s.nodeID != "" { + return s.nodeID } nodeID, err := s.holder.loadNodeID() if err != nil { s.logger.Printf("loading NodeID: %v", err) - return s.NodeID + return s.nodeID } return nodeID } -type pilosaAddr URI - -func (p pilosaAddr) String() string { - uri := URI(p) - return uri.HostPort() - -} - -func (pilosaAddr) Network() string { - return "tcp" -} - -// Addr returns the address of the listener. -func (s *Server) Addr() net.Addr { - return pilosaAddr(s.URI) -} - func (s *Server) monitorAntiEntropy() { ticker := time.NewTicker(s.antiEntropyInterval) defer ticker.Stop() @@ -396,8 +391,8 @@ func (s *Server) monitorAntiEntropy() { // Initialize syncer with local holder and remote client. var syncer HolderSyncer syncer.Holder = s.holder - syncer.Node = s.Cluster.Node - syncer.Cluster = s.Cluster + syncer.Node = s.cluster.Node + syncer.Cluster = s.cluster syncer.Closing = s.closing syncer.Stats = s.holder.Stats.WithTags("HolderSyncer") @@ -467,33 +462,33 @@ func (s *Server) ReceiveMessage(pb proto.Message) error { return err } case *internal.ClusterStatus: - err := s.Cluster.mergeClusterStatus(obj) + err := s.cluster.mergeClusterStatus(obj) if err != nil { return err } case *internal.ResizeInstruction: - err := s.Cluster.followResizeInstruction(obj) + err := s.cluster.followResizeInstruction(obj) if err != nil { return err } case *internal.ResizeInstructionComplete: - err := s.Cluster.markResizeInstructionComplete(obj) + err := s.cluster.markResizeInstructionComplete(obj) if err != nil { return err } case *internal.SetCoordinatorMessage: - s.Cluster.setCoordinator(DecodeNode(obj.New)) + s.cluster.setCoordinator(DecodeNode(obj.New)) case *internal.UpdateCoordinatorMessage: - s.Cluster.updateCoordinator(DecodeNode(obj.New)) + s.cluster.updateCoordinator(DecodeNode(obj.New)) case *internal.NodeStateMessage: - err := s.Cluster.receiveNodeState(obj.NodeID, obj.State) + err := s.cluster.receiveNodeState(obj.NodeID, obj.State) if err != nil { return err } case *internal.RecalculateCaches: s.holder.RecalculateCaches() case *internal.NodeEventMessage: - s.Cluster.ReceiveEvent(DecodeNodeEvent(obj)) + s.cluster.ReceiveEvent(DecodeNodeEvent(obj)) } return nil @@ -502,7 +497,7 @@ func (s *Server) ReceiveMessage(pb proto.Message) error { // SendSync represents an implementation of Broadcaster. func (s *Server) SendSync(pb proto.Message) error { var eg errgroup.Group - for _, node := range s.Cluster.Nodes { + for _, node := range s.cluster.Nodes { node := node s.logger.Printf("SendSync to: %s", node.URI) // Don't forward the message to ourselves. @@ -529,6 +524,12 @@ func (s *Server) SendTo(to *Node, pb proto.Message) error { return s.defaultClient.SendMessage(context.Background(), &to.URI, pb) } +// Node returns the pilosa.Node object. It is used by membership protocols to +// get this node's name(ID), location(URI), and coordinator status. +func (s *Server) Node() *Node { + return s.cluster.Node +} + // Server implements StatusHandler. // LocalStatus is used to periodically sync information // between nodes. Under normal conditions, nodes should @@ -540,7 +541,7 @@ func (s *Server) SendTo(to *Node, pb proto.Message) error { // - Schema // In a gossip implementation, memberlist.Delegate.LocalState() uses this. func (s *Server) LocalStatus() (proto.Message, error) { - if s.Cluster == nil { + if s.cluster == nil { return nil, errors.New("Server.Cluster is nil") } if s.holder == nil { @@ -548,7 +549,7 @@ func (s *Server) LocalStatus() (proto.Message, error) { } ns := internal.NodeStatus{ - Node: EncodeNode(s.Cluster.Node), + Node: EncodeNode(s.cluster.Node), MaxSlices: s.holder.EncodeMaxSlices(), Schema: s.holder.EncodeSchema(), } @@ -558,13 +559,13 @@ func (s *Server) LocalStatus() (proto.Message, error) { // ClusterStatus returns the ClusterState and NodeSet for the cluster. func (s *Server) ClusterStatus() (proto.Message, error) { - return s.Cluster.Status(), nil + return s.cluster.Status(), nil } // HandleRemoteStatus receives incoming NodeStatus from remote nodes. func (s *Server) HandleRemoteStatus(pb proto.Message) error { // Ignore NodeStatus messages until the cluster is in a Normal state. - if s.Cluster.State() != ClusterStateNormal { + if s.cluster.State() != ClusterStateNormal { return nil } @@ -583,7 +584,7 @@ func (s *Server) HandleRemoteStatus(pb proto.Message) error { func (s *Server) mergeRemoteStatus(ns *internal.NodeStatus) error { // Ignore status updates from self. - if s.NodeID == DecodeNode(ns.Node).ID { + if s.nodeID == DecodeNode(ns.Node).ID { return nil } @@ -624,11 +625,11 @@ func (s *Server) monitorDiagnostics() { s.diagnostics.Logger = s.logger s.diagnostics.SetVersion(Version) s.diagnostics.Set("Host", s.URI.host) - s.diagnostics.Set("Cluster", strings.Join(s.Cluster.nodeIDs(), ",")) - s.diagnostics.Set("NumNodes", len(s.Cluster.Nodes)) + s.diagnostics.Set("Cluster", strings.Join(s.cluster.nodeIDs(), ",")) + s.diagnostics.Set("NumNodes", len(s.cluster.Nodes)) s.diagnostics.Set("NumCPU", runtime.NumCPU()) - s.diagnostics.Set("NodeID", s.NodeID) - s.diagnostics.Set("ClusterID", s.Cluster.ID) + s.diagnostics.Set("NodeID", s.nodeID) + s.diagnostics.Set("ClusterID", s.cluster.id) s.diagnostics.EnrichWithOSInfo() // Flush the diagnostics metrics at startup, then on each tick interval @@ -708,7 +709,7 @@ func (s *Server) monitorRuntime() { // ReceiveEvent implements the EventHandler interface. func (s *Server) ReceiveEvent(e *NodeEvent) error { - return s.Cluster.ReceiveEvent(e) + return s.cluster.ReceiveEvent(e) } // countOpenFiles on operating systems that support lsof. diff --git a/server/cluster_test.go b/server/cluster_test.go index b9a1917a9..825911ace 100644 --- a/server/cluster_test.go +++ b/server/cluster_test.go @@ -24,10 +24,9 @@ import ( "testing" "time" - "golang.org/x/sync/errgroup" - "github.com/pilosa/pilosa" "github.com/pilosa/pilosa/test" + "golang.org/x/sync/errgroup" ) // Ensure program can send/receive broadcast messages. @@ -37,11 +36,6 @@ func TestMain_SendReceiveMessage(t *testing.T) { defer m0.Close() defer m1.Close() - m0.Server.Cluster.SetState(pilosa.ClusterStateNormal) - m1.Server.Cluster.SetState(pilosa.ClusterStateNormal) - - //////////////////////////////////////////////////////////////////////////////////////////////////////////////////// - // Expected indexes and Fields expected := map[string][]string{ "i": []string{"f"}, @@ -125,83 +119,41 @@ func TestClusterResize_EmptyNode(t *testing.T) { m0 := test.MustRunMain() defer m0.Close() - if m0.Server.Cluster.State() != pilosa.ClusterStateNormal { - t.Fatalf("unexpected cluster state: %s", m0.Server.Cluster.State()) + if m0.API.State() != pilosa.ClusterStateNormal { + t.Fatalf("unexpected cluster state: %s", m0.API.State()) } } // Ensure that a cluster of empty nodes comes up in a NORMAL state. func TestClusterResize_EmptyNodes(t *testing.T) { - // Configure node0 - m0 := test.NewMainWithCluster(true) - defer m0.Close() + clus := test.MustRunMainWithCluster(t, 2) + defer clus[0].Close() + defer clus[1].Close() - gossipHost := "localhost" - gossipPort := 0 - seed, err := m0.RunWithTransport(gossipHost, gossipPort, []string{}) - if err != nil { - t.Fatal(err) - } - - // Configure node1 - m1 := test.NewMainWithCluster(false) - defer m1.Close() - - seed, err = m1.RunWithTransport(gossipHost, gossipPort, []string{seed}) - if err != nil { - t.Fatal(err) - } - - if m0.Server.Cluster.State() != pilosa.ClusterStateNormal { - t.Fatalf("unexpected node0 cluster state: %s", m0.Server.Cluster.State()) - } else if m1.Server.Cluster.State() != pilosa.ClusterStateNormal { - t.Fatalf("unexpected node1 cluster state: %s", m1.Server.Cluster.State()) + if clus[0].API.State() != pilosa.ClusterStateNormal { + t.Fatalf("unexpected node0 cluster state: %s", clus[0].API.State()) + } else if clus[1].API.State() != pilosa.ClusterStateNormal { + t.Fatalf("unexpected node1 cluster state: %s", clus[1].API.State()) } } // Ensure that adding a node correctly resizes the cluster. func TestClusterResize_AddNode(t *testing.T) { t.Run("NoData", func(t *testing.T) { - // Configure node0 - m0 := test.NewMainWithCluster(true) - defer m0.Close() + clus := test.MustRunMainWithCluster(t, 2) - seed, err := m0.RunWithTransport("localhost", 0, []string{}) - if err != nil { - t.Fatal(err) - } - - // Configure node1 - m1 := test.NewMainWithCluster(false) - defer m1.Close() - - var eg errgroup.Group - eg.Go(func() error { - _, err = m1.RunWithTransport("localhost", 0, []string{seed}) - if err != nil { - return err - } - return nil - }) - if err := eg.Wait(); err != nil { - t.Fatal(err) - } - - if !checkClusterState(m0.Server.Cluster, pilosa.ClusterStateNormal, 1000) { - t.Fatalf("unexpected node0 cluster state: %s", m0.Server.Cluster.State()) - } else if !checkClusterState(m1.Server.Cluster, pilosa.ClusterStateNormal, 1000) { - t.Fatalf("unexpected node1 cluster state: %s", m1.Server.Cluster.State()) + if !checkClusterState(clus[0], pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node0 cluster state: %s", clus[0].API.State()) + } else if !checkClusterState(clus[1], pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node1 cluster state: %s", clus[1].API.State()) } }) t.Run("WithIndex", func(t *testing.T) { // Configure node0 - m0 := test.NewMainWithCluster(true) + m0 := test.MustRunMainWithCluster(t, 1)[0] defer m0.Close() - seed, err := m0.RunWithTransport("localhost", 0, []string{}) - if err != nil { - t.Fatal(err) - } + seed := m0.GossipAddress() // Create a client for each node. client0 := m0.Client() @@ -215,40 +167,29 @@ func TestClusterResize_AddNode(t *testing.T) { // Configure node1 m1 := test.NewMainWithCluster(false) + m1.Config.Gossip.Port = "0" + m1.Config.Gossip.Seeds = []string{seed} + err := m1.Start() + if err != nil { + t.Fatalf("starting second main: %v", err) + } defer m1.Close() - var eg errgroup.Group - eg.Go(func() error { - _, err = m1.RunWithTransport("localhost", 0, []string{seed}) - if err != nil { - return err - } - return nil - }) - if err := eg.Wait(); err != nil { - t.Fatal(err) - } - - if !checkClusterState(m0.Server.Cluster, pilosa.ClusterStateNormal, 1000) { - t.Fatalf("unexpected node0 cluster state: %s", m0.Server.Cluster.State()) - } else if !checkClusterState(m1.Server.Cluster, pilosa.ClusterStateNormal, 1000) { - t.Fatalf("unexpected node1 cluster state: %s", m1.Server.Cluster.State()) + if !checkClusterState(m0, pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node0 cluster state: %s", m0.API.State()) + } else if !checkClusterState(m1, pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node1 cluster state: %s", m1.API.State()) } }) t.Run("ContinuousSlices", func(t *testing.T) { - // Configure node0 - m0 := test.NewMainWithCluster(true) + m0 := test.MustRunMainWithCluster(t, 1)[0] defer m0.Close() - seed, err := m0.RunWithTransport("localhost", 0, []string{}) - if err != nil { - t.Fatal(err) - } + seed := m0.GossipAddress() // Create a client for each node. client0 := m0.Client() - //client1 := m1.Client() // Create indexes and fields on one node. if err := client0.CreateIndex(context.Background(), "i", pilosa.IndexOptions{}); err != nil && err != pilosa.ErrIndexExists { @@ -267,40 +208,29 @@ func TestClusterResize_AddNode(t *testing.T) { // Configure node1 m1 := test.NewMainWithCluster(false) + m1.Config.Gossip.Port = "0" + m1.Config.Gossip.Seeds = []string{seed} + err := m1.Start() + if err != nil { + t.Fatalf("starting second main: %v", err) + } defer m1.Close() - var eg errgroup.Group - eg.Go(func() error { - _, err = m1.RunWithTransport("localhost", 0, []string{seed}) - if err != nil { - return err - } - return nil - }) - if err := eg.Wait(); err != nil { - t.Fatal(err) - } - - if !checkClusterState(m0.Server.Cluster, pilosa.ClusterStateNormal, 1000) { - t.Fatalf("unexpected node0 cluster state: %s", m0.Server.Cluster.State()) - } else if !checkClusterState(m1.Server.Cluster, pilosa.ClusterStateNormal, 1000) { - t.Fatalf("unexpected node1 cluster state: %s", m1.Server.Cluster.State()) + if !checkClusterState(m0, pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node0 cluster state: %s", m0.API.State()) + } else if !checkClusterState(m1, pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node1 cluster state: %s", m1.API.State()) } }) t.Run("SkippedSlice", func(t *testing.T) { - // Configure node0 - m0 := test.NewMainWithCluster(true) + m0 := test.MustRunMainWithCluster(t, 1)[0] defer m0.Close() - seed, err := m0.RunWithTransport("localhost", 0, []string{}) - if err != nil { - t.Fatal(err) - } + seed := m0.GossipAddress() // Create a client for each node. client0 := m0.Client() - //client1 := m1.Client() // Create indexes and fields on one node. if err := client0.CreateIndex(context.Background(), "i", pilosa.IndexOptions{}); err != nil && err != pilosa.ErrIndexExists { @@ -319,24 +249,18 @@ func TestClusterResize_AddNode(t *testing.T) { // Configure node1 m1 := test.NewMainWithCluster(false) + m1.Config.Gossip.Port = "0" + m1.Config.Gossip.Seeds = []string{seed} + err := m1.Start() + if err != nil { + t.Fatalf("starting second main: %v", err) + } defer m1.Close() - var eg errgroup.Group - eg.Go(func() error { - _, err = m1.RunWithTransport("localhost", 0, []string{seed}) - if err != nil { - return err - } - return nil - }) - if err := eg.Wait(); err != nil { - t.Fatal(err) - } - - if !checkClusterState(m0.Server.Cluster, pilosa.ClusterStateNormal, 1000) { - t.Fatalf("unexpected node0 cluster state: %s", m0.Server.Cluster.State()) - } else if !checkClusterState(m1.Server.Cluster, pilosa.ClusterStateNormal, 1000) { - t.Fatalf("unexpected node1 cluster state: %s", m1.Server.Cluster.State()) + if !checkClusterState(m0, pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node0 cluster state: %s", m0.API.State()) + } else if !checkClusterState(m1, pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node1 cluster state: %s", m1.API.State()) } }) } @@ -345,37 +269,37 @@ func TestClusterResize_AddNode(t *testing.T) { func TestCluster_GossipMembership(t *testing.T) { t.Run("Node0Down", func(t *testing.T) { // Configure node0 - m0 := test.NewMainWithCluster(true) + m0 := test.MustRunMainWithCluster(t, 1)[0] defer m0.Close() - seed, err := m0.RunWithTransport("localhost", 0, []string{}) - if err != nil { - t.Fatal(err) - } + seed := m0.GossipAddress() + + var eg errgroup.Group // Configure node1 m1 := test.NewMainWithCluster(false) defer m1.Close() - - var eg errgroup.Group eg.Go(func() error { + m1.Config.Gossip.Port = "0" // Pass invalid seed as first in list - _, err := m1.RunWithTransport("localhost", 0, []string{"http://localhost:8765", seed}) + m1.Config.Gossip.Seeds = []string{"http://localhost:8765", seed} + err := m1.Start() if err != nil { - return err + t.Fatalf("starting second main: %v", err) } return nil }) - // Configure node2 + // Configure node1 m2 := test.NewMainWithCluster(false) defer m2.Close() - eg.Go(func() error { - // Pass invalid seed as last in list - _, err := m2.RunWithTransport("localhost", 0, []string{seed, "http://localhost:8765"}) + m2.Config.Gossip.Port = "0" + // Pass invalid seed as first in list + m2.Config.Gossip.Seeds = []string{seed, "http://localhost:8765"} + err := m2.Start() if err != nil { - return err + t.Fatalf("starting second main: %v", err) } return nil }) @@ -384,15 +308,15 @@ func TestCluster_GossipMembership(t *testing.T) { t.Fatal(err) } - if !checkClusterState(m0.Server.Cluster, pilosa.ClusterStateNormal, 1000) { - t.Fatalf("unexpected node0 cluster state: %s", m0.Server.Cluster.State()) - } else if !checkClusterState(m1.Server.Cluster, pilosa.ClusterStateNormal, 1000) { - t.Fatalf("unexpected node1 cluster state: %s", m1.Server.Cluster.State()) - } else if !checkClusterState(m2.Server.Cluster, pilosa.ClusterStateNormal, 1000) { - t.Fatalf("unexpected node2 cluster state: %s", m2.Server.Cluster.State()) + if !checkClusterState(m0, pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node0 cluster state: %s", m0.API.State()) + } else if !checkClusterState(m1, pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node1 cluster state: %s", m1.API.State()) + } else if !checkClusterState(m2, pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node2 cluster state: %s", m2.API.State()) } - numNodes := len(m0.Server.Cluster.Status().Nodes) + numNodes := len(m0.API.Hosts(context.Background())) if numNodes != 3 { t.Fatalf("Expected 3 nodes, got %d", numNodes) } @@ -486,9 +410,9 @@ func TestClusterResize_RemoveNode(t *testing.T) { // checkClusterState polls a given cluster for its state until it // receives a matching state. It polls up to n times before returning. -func checkClusterState(c *pilosa.Cluster, state string, n int) bool { +func checkClusterState(m *test.Main, state string, n int) bool { for i := 0; i < n; i++ { - if c.State() == state { + if m.API.State() == state { return true } time.Sleep(10 * time.Millisecond) diff --git a/server/server.go b/server/server.go index a6bc1d995..164dc94e3 100644 --- a/server/server.go +++ b/server/server.go @@ -60,7 +60,7 @@ type Command struct { Config *Config // Gossip transport - GossipTransport *gossip.Transport + gossipTransport *gossip.Transport // Standard input/output *pilosa.CmdIO @@ -247,6 +247,12 @@ func (m *Command) SetupServer() error { primaryTranslateStore = http.NewTranslateStore(m.Config.Translation.PrimaryURL) } + // Set Coordinator. + coordinatorOpt := pilosa.OptServerIsCoordinator(false) + if m.Config.Cluster.Coordinator || len(m.Config.Gossip.Seeds) == 0 { + coordinatorOpt = pilosa.OptServerIsCoordinator(true) + } + serverOptions := []pilosa.ServerOption{ pilosa.OptServerAntiEntropyInterval(time.Duration(m.Config.AntiEntropy.Interval)), pilosa.OptServerLongQueryTime(time.Duration(m.Config.Cluster.LongQueryTime)), @@ -265,6 +271,7 @@ func (m *Command) SetupServer() error { pilosa.OptServerInternalClient(http.NewInternalClientFromURI(uri, c)), pilosa.OptServerPrimaryTranslateStore(primaryTranslateStore), pilosa.OptServerClusterDisabled(m.Config.Cluster.Disabled, m.Config.Cluster.Hosts), + coordinatorOpt, } serverOptions = append(serverOptions, m.serverOptions...) @@ -303,35 +310,28 @@ func (m *Command) SetupNetworking() error { // get the host portion of addr to use for binding gossipHost := m.Server.URI.Host() - var transport *gossip.Transport - if m.GossipTransport != nil { - transport = m.GossipTransport - } else { - transport, err = gossip.NewTransport(gossipHost, gossipPort, m.logger.Logger()) - if err != nil { - return errors.Wrap(err, "getting transport") - } - } - - // Set Coordinator. - if m.Config.Cluster.Coordinator || len(m.Config.Gossip.Seeds) == 0 { - m.Server.Cluster.Coordinator = m.Server.NodeID - m.Server.Cluster.Node.IsCoordinator = true + m.gossipTransport, err = gossip.NewTransport(gossipHost, gossipPort, m.logger.Logger()) + if err != nil { + return errors.Wrap(err, "getting transport") } gossipMemberSet, err := gossip.NewGossipMemberSet( - m.Server.NodeID, - m.Server.URI.Host(), m.Config.Gossip, m.Server, gossip.WithLogger(m.logger.Logger()), - gossip.WithTransport(transport), + gossip.WithTransport(m.gossipTransport), ) if err != nil { return errors.Wrap(err, "getting memberset") } - m.Server.Cluster.MemberSet = gossipMemberSet - return nil + return errors.Wrap(gossipMemberSet.Open(), "opening gossip memberset") +} + +// GossipTransport allows a caller to return the gossip transport created when +// setting up the GossipMemberSet. This is useful if one needs to determine the +// allocated ephemeral port programmatically. (usually used in tests) +func (m *Command) GossipTransport() *gossip.Transport { + return m.gossipTransport } // Close shuts down the server. diff --git a/test/pilosa.go b/test/pilosa.go index 47b132508..e201de283 100644 --- a/test/pilosa.go +++ b/test/pilosa.go @@ -19,14 +19,12 @@ import ( "fmt" "io" "io/ioutil" - "log" gohttp "net/http" "os" "strings" "testing" "time" - "github.com/pilosa/pilosa/gossip" "github.com/pilosa/pilosa/http" "github.com/pilosa/pilosa/server" "github.com/pilosa/pilosa/toml" @@ -59,6 +57,13 @@ func OptAllowedOrigins(origins []string) server.CommandOption { } } +// GossipAddress returns the address on which gossip is listening after a Main +// has been setup. Useful to pass as a seed to other nodes when creating and +// testing clusters. +func (m *Main) GossipAddress() string { + return m.GossipTransport().URI.String() +} + // NewMain returns a new instance of Main with a temporary data directory and random port. func NewMain(opts ...server.CommandOption) *Main { path, err := ioutil.TempDir("", "pilosa-") @@ -116,25 +121,20 @@ func runMainWithCluster(size int, opts ...[]server.CommandOption) ([]*Main, erro } mains := make([]*Main, size) - - gossipHost := "localhost" - gossipPort := 0 - var err error var gossipSeeds = make([]string, size) - for i := 0; i < size; i++ { var commandOpts []server.CommandOption if len(opts) > 0 { commandOpts = opts[i%len(opts)] } m := NewMainWithCluster(i == 0, commandOpts...) - m.Config.Cluster.Disabled = false + m.Config.Gossip.Port = "0" + m.Config.Gossip.Seeds = gossipSeeds[:i] - gossipSeeds[i], err = m.RunWithTransport(gossipHost, gossipPort, gossipSeeds[:i]) - if err != nil { - return nil, errors.Wrap(err, "RunWithTransport") + if err := m.Start(); err != nil { + return nil, errors.Wrapf(err, "Starting server %d", i) } - + gossipSeeds[i] = m.GossipTransport().URI.String() mains[i] = m } @@ -179,71 +179,8 @@ func (m *Main) Reopen() error { return nil } -// RunWithTransport runs Main and returns the dynamically allocated gossip port. -func (m *Main) RunWithTransport(host string, bindPort int, joinSeeds []string) (seed string, err error) { - defer close(m.Started) - - /* - TEST: - - SetupServer (just static settings from config) - - OpenListener (sets Server.Name to use in gossip) - - NewTransport (gossip) - - SetupNetworking (does the gossip or static stuff) - uses Server.Name - - Open server - - PRODUCTION: - - SetupServer (just static settings from config) - - SetupNetworking (does the gossip or static stuff) - calls NewTransport - - Open server - calls OpenListener - */ - - // SetupServer - err = m.SetupServer() - if err != nil { - return seed, err - } - - // Open gossip transport to use in SetupServer. - transport, err := gossip.NewTransport(host, bindPort, nil) - if err != nil { - return seed, err - } - m.GossipTransport = transport - - if len(joinSeeds) != 0 { - m.Config.Gossip.Seeds = joinSeeds - } else { - m.Config.Gossip.Seeds = []string{transport.URI.String()} - } - - seed = transport.URI.String() - - // SetupNetworking - err = m.SetupNetworking() - if err != nil { - return seed, err - } - - m.Server.Cluster.Static = false - - go func() { - err := m.Handler.Serve() - if err != nil { - log.Printf("Handler serve error: %v", err) - } - }() - - // Initialize server. - err = m.Server.Open() - if err != nil { - return seed, err - } - - return seed, nil -} - // URL returns the base URL string for accessing the running program. -func (m *Main) URL() string { return "http://" + m.Server.Addr().String() } +func (m *Main) URL() string { return m.Server.URI.String() } // Client returns a client to connect to the program. func (m *Main) Client() *http.InternalClient { diff --git a/test/pilosa_test.go b/test/pilosa_test.go index 0a1acc551..2ba7504c8 100644 --- a/test/pilosa_test.go +++ b/test/pilosa_test.go @@ -15,6 +15,7 @@ package test_test import ( + "context" "encoding/json" "net/http" "strings" @@ -27,15 +28,16 @@ import ( func TestNewCluster(t *testing.T) { numNodes := 3 cluster := test.MustRunMainWithCluster(t, numNodes) - coordinator := cluster[0].Server.Cluster.Coordinator + + coordinator := getCoordinator(cluster[0]) for i := 1; i < numNodes; i++ { - if coordi := cluster[i].Server.Cluster.Coordinator; coordi != coordinator { + if coordi := getCoordinator(cluster[i]); coordi != coordinator { t.Fatalf("node %d does not have the same coordinator as node 0. '%v' and '%v' respectively", i, coordi, coordinator) } } req, err := http.NewRequest( "GET", - "http://"+cluster[0].Server.Addr().String()+"/status", + cluster[0].URL()+"/status", strings.NewReader(""), ) @@ -75,3 +77,13 @@ func TestNewCluster(t *testing.T) { t.Fatalf("cluster state should be %s but is %s", pilosa.ClusterStateNormal, body.State) } } + +func getCoordinator(m *test.Main) string { + hosts := m.API.Hosts(context.Background()) + for _, host := range hosts { + if host.IsCoordinator { + return host.ID + } + } + panic("no coordinator in cluster") +} diff --git a/utils_internal_test.go b/utils_internal_test.go index ea98875c4..5e19a83b9 100644 --- a/utils_internal_test.go +++ b/utils_internal_test.go @@ -97,7 +97,7 @@ type commonClusterSettings struct { func (t *ClusterCluster) CreateIndex(name string) error { for _, c := range t.Clusters { - if _, err := c.Holder.CreateIndexIfNotExists(name, IndexOptions{}); err != nil { + if _, err := c.holder.CreateIndexIfNotExists(name, IndexOptions{}); err != nil { return err } } @@ -106,7 +106,7 @@ func (t *ClusterCluster) CreateIndex(name string) error { func (t *ClusterCluster) CreateField(index, field string, opt FieldOptions) error { for _, c := range t.Clusters { - idx, err := c.Holder.CreateIndexIfNotExists(index, IndexOptions{}) + idx, err := c.holder.CreateIndexIfNotExists(index, IndexOptions{}) if err != nil { return err } @@ -128,7 +128,7 @@ func (t *ClusterCluster) SetBit(index, field string, rowID, colID uint64, x *tim if c == nil { continue } - f := c.Holder.Field(index, field) + f := c.holder.Field(index, field) if f == nil { return fmt.Errorf("index/field does not exist: %s/%s", index, field) } @@ -227,11 +227,10 @@ func (t *ClusterCluster) addCluster(i int, saveTopology bool) (*Cluster, error) c.Hasher = NewTestModHasher() c.Path = path c.Topology = NewTopology() - c.Holder = h - c.MemberSet = NewStaticMemberSet(c.Nodes) + c.holder = h c.Node = node c.Coordinator = t.common.Nodes[0].ID // the first node is the coordinator - c.Broadcaster = t + c.broadcaster = t // add nodes if saveTopology { @@ -276,7 +275,7 @@ func (t *ClusterCluster) Open() error { if err := c.open(); err != nil { return err } - if err := c.Holder.Open(); err != nil { + if err := c.holder.Open(); err != nil { return err } if err := c.setNodeState(NodeStateReady); err != nil { @@ -361,7 +360,7 @@ func (t *ClusterCluster) FollowResizeInstruction(instr *internal.ResizeInstructi destCluster := t.clusterByID(instrNode.ID) // Sync the schema received in the resize instruction. - if err := destCluster.Holder.ApplySchema(instr.Schema); err != nil { + if err := destCluster.holder.ApplySchema(instr.Schema); err != nil { return err } @@ -369,11 +368,11 @@ func (t *ClusterCluster) FollowResizeInstruction(instr *internal.ResizeInstructi srcNode := DecodeNode(src.Node) srcCluster := t.clusterByID(srcNode.ID) - srcFragment := srcCluster.Holder.Fragment(src.Index, src.Field, src.View, src.Slice) - destFragment := destCluster.Holder.Fragment(src.Index, src.Field, src.View, src.Slice) + srcFragment := srcCluster.holder.Fragment(src.Index, src.Field, src.View, src.Slice) + destFragment := destCluster.holder.Fragment(src.Index, src.Field, src.View, src.Slice) if destFragment == nil { // Create fragment on destination if it doesn't exist. - f := destCluster.Holder.Field(src.Index, src.Field) + f := destCluster.holder.Field(src.Index, src.Field) v := f.View(src.View) var err error destFragment, err = v.CreateFragmentIfNotExists(src.Slice)