From bcd15ef271bcea74eaf01a7e00a82acbcde5c96c Mon Sep 17 00:00:00 2001 From: Ben Johnson Date: Mon, 9 Jul 2018 18:31:27 +0100 Subject: [PATCH 1/2] peg: Allow single quoted key 'col'. --- pql/pql.peg | 1 + pql/pql.peg.go | 530 ++++++++++++++++++++++++--------------------- pql/pqlpeg_test.go | 8 + 3 files changed, 290 insertions(+), 249 deletions(-) diff --git a/pql/pql.peg b/pql/pql.peg index 141d8f265..f55267e1f 100644 --- a/pql/pql.peg +++ b/pql/pql.peg @@ -57,6 +57,7 @@ posfield <- { p.addPosStr("_field", buffer[begin:end]) } uint <- [1-9] [0-9]* / '0' uintrow <- {p.addPosNum("_row", buffer[begin:end])} col <- ( {p.addPosNum("_col", buffer[begin:end])} + / '\'' '\'' {p.addPosStr("_col", buffer[begin:end])} / '"' '"' {p.addPosStr("_col", buffer[begin:end])} ) diff --git a/pql/pql.peg.go b/pql/pql.peg.go index 2516def6c..7092b8f1f 100644 --- a/pql/pql.peg.go +++ b/pql/pql.peg.go @@ -94,6 +94,7 @@ const ( ruleAction41 ruleAction42 ruleAction43 + ruleAction44 ) var rul3s = [...]string{ @@ -176,6 +177,7 @@ var rul3s = [...]string{ "Action41", "Action42", "Action43", + "Action44", } type token32 struct { @@ -292,7 +294,7 @@ type PQL struct { Buffer string buffer []rune - rules [79]func() bool + rules [80]func() bool parse func(rule ...int) error reset func() Pretty bool @@ -471,6 +473,8 @@ func (p *PQL) Execute() { case ruleAction42: p.addPosStr("_col", buffer[begin:end]) case ruleAction43: + p.addPosStr("_col", buffer[begin:end]) + case ruleAction44: p.addPosStr("_timestamp", buffer[begin:end]) } @@ -632,7 +636,7 @@ func (p *PQL) Init() { add(rulePegText, position13) } { - add(ruleAction43, position) + add(ruleAction44, position) } add(ruletimestamp, position12) } @@ -1912,95 +1916,8 @@ func (p *PQL) Init() { position++ { position183 := position - { - position184 := position - l185: - { - position186, tokenIndex186 := position, tokenIndex - { - position187, tokenIndex187 := position, tokenIndex - { - position189, tokenIndex189 := position, tokenIndex - { - position190, tokenIndex190 := position, tokenIndex - if buffer[position] != rune('\'') { - goto l191 - } - position++ - goto l190 - l191: - position, tokenIndex = position190, tokenIndex190 - if buffer[position] != rune('\\') { - goto l192 - } - position++ - goto l190 - l192: - position, tokenIndex = position190, tokenIndex190 - if buffer[position] != rune('\n') { - goto l189 - } - position++ - } - l190: - goto l188 - l189: - position, tokenIndex = position189, tokenIndex189 - } - if !matchDot() { - goto l188 - } - goto l187 - l188: - position, tokenIndex = position187, tokenIndex187 - if buffer[position] != rune('\\') { - goto l193 - } - position++ - if buffer[position] != rune('n') { - goto l193 - } - position++ - goto l187 - l193: - position, tokenIndex = position187, tokenIndex187 - if buffer[position] != rune('\\') { - goto l194 - } - position++ - if buffer[position] != rune('"') { - goto l194 - } - position++ - goto l187 - l194: - position, tokenIndex = position187, tokenIndex187 - if buffer[position] != rune('\\') { - goto l195 - } - position++ - if buffer[position] != rune('\'') { - goto l195 - } - position++ - goto l187 - l195: - position, tokenIndex = position187, tokenIndex187 - if buffer[position] != rune('\\') { - goto l186 - } - position++ - if buffer[position] != rune('\\') { - goto l186 - } - position++ - } - l187: - goto l185 - l186: - position, tokenIndex = position186, tokenIndex186 - } - add(rulesinglequotedstring, position184) + if !_rules[rulesinglequotedstring]() { + goto l127 } add(rulePegText, position183) } @@ -2023,99 +1940,191 @@ func (p *PQL) Init() { /* 14 doublequotedstring <- <((!('"' / '\\' / '\n') .) / ('\\' 'n') / ('\\' '"') / ('\\' '\'') / ('\\' '\\'))*> */ func() bool { { - position198 := position - l199: + position186 := position + l187: { - position200, tokenIndex200 := position, tokenIndex + position188, tokenIndex188 := position, tokenIndex { - position201, tokenIndex201 := position, tokenIndex + position189, tokenIndex189 := position, tokenIndex { - position203, tokenIndex203 := position, tokenIndex + position191, tokenIndex191 := position, tokenIndex { - position204, tokenIndex204 := position, tokenIndex + position192, tokenIndex192 := position, tokenIndex if buffer[position] != rune('"') { - goto l205 + goto l193 } position++ - goto l204 - l205: - position, tokenIndex = position204, tokenIndex204 + goto l192 + l193: + position, tokenIndex = position192, tokenIndex192 if buffer[position] != rune('\\') { - goto l206 + goto l194 } position++ - goto l204 - l206: - position, tokenIndex = position204, tokenIndex204 + goto l192 + l194: + position, tokenIndex = position192, tokenIndex192 if buffer[position] != rune('\n') { - goto l203 + goto l191 } position++ } - l204: - goto l202 - l203: - position, tokenIndex = position203, tokenIndex203 + l192: + goto l190 + l191: + position, tokenIndex = position191, tokenIndex191 } if !matchDot() { - goto l202 + goto l190 } - goto l201 - l202: - position, tokenIndex = position201, tokenIndex201 + goto l189 + l190: + position, tokenIndex = position189, tokenIndex189 if buffer[position] != rune('\\') { - goto l207 + goto l195 } position++ if buffer[position] != rune('n') { - goto l207 + goto l195 } position++ - goto l201 - l207: - position, tokenIndex = position201, tokenIndex201 + goto l189 + l195: + position, tokenIndex = position189, tokenIndex189 if buffer[position] != rune('\\') { - goto l208 + goto l196 } position++ if buffer[position] != rune('"') { - goto l208 + goto l196 } position++ - goto l201 - l208: - position, tokenIndex = position201, tokenIndex201 + goto l189 + l196: + position, tokenIndex = position189, tokenIndex189 if buffer[position] != rune('\\') { - goto l209 + goto l197 } position++ if buffer[position] != rune('\'') { - goto l209 + goto l197 } position++ - goto l201 - l209: - position, tokenIndex = position201, tokenIndex201 + goto l189 + l197: + position, tokenIndex = position189, tokenIndex189 if buffer[position] != rune('\\') { - goto l200 + goto l188 } position++ if buffer[position] != rune('\\') { - goto l200 + goto l188 } position++ } - l201: - goto l199 - l200: - position, tokenIndex = position200, tokenIndex200 + l189: + goto l187 + l188: + position, tokenIndex = position188, tokenIndex188 } - add(ruledoublequotedstring, position198) + add(ruledoublequotedstring, position186) } return true }, /* 15 singlequotedstring <- <((!('\'' / '\\' / '\n') .) / ('\\' 'n') / ('\\' '"') / ('\\' '\'') / ('\\' '\\'))*> */ - nil, + func() bool { + { + position199 := position + l200: + { + position201, tokenIndex201 := position, tokenIndex + { + position202, tokenIndex202 := position, tokenIndex + { + position204, tokenIndex204 := position, tokenIndex + { + position205, tokenIndex205 := position, tokenIndex + if buffer[position] != rune('\'') { + goto l206 + } + position++ + goto l205 + l206: + position, tokenIndex = position205, tokenIndex205 + if buffer[position] != rune('\\') { + goto l207 + } + position++ + goto l205 + l207: + position, tokenIndex = position205, tokenIndex205 + if buffer[position] != rune('\n') { + goto l204 + } + position++ + } + l205: + goto l203 + l204: + position, tokenIndex = position204, tokenIndex204 + } + if !matchDot() { + goto l203 + } + goto l202 + l203: + position, tokenIndex = position202, tokenIndex202 + if buffer[position] != rune('\\') { + goto l208 + } + position++ + if buffer[position] != rune('n') { + goto l208 + } + position++ + goto l202 + l208: + position, tokenIndex = position202, tokenIndex202 + if buffer[position] != rune('\\') { + goto l209 + } + position++ + if buffer[position] != rune('"') { + goto l209 + } + position++ + goto l202 + l209: + position, tokenIndex = position202, tokenIndex202 + if buffer[position] != rune('\\') { + goto l210 + } + position++ + if buffer[position] != rune('\'') { + goto l210 + } + position++ + goto l202 + l210: + position, tokenIndex = position202, tokenIndex202 + if buffer[position] != rune('\\') { + goto l201 + } + position++ + if buffer[position] != rune('\\') { + goto l201 + } + position++ + } + l202: + goto l200 + l201: + position, tokenIndex = position201, tokenIndex201 + } + add(rulesinglequotedstring, position199) + } + return true + }, /* 16 fieldExpr <- <(([a-z] / [A-Z]) ([a-z] / [A-Z] / [0-9] / '_' / '-')*)> */ func() bool { position211, tokenIndex211 := position, tokenIndex @@ -2438,7 +2447,7 @@ func (p *PQL) Init() { }, /* 21 uintrow <- <( Action40)> */ nil, - /* 22 col <- <(( Action41) / ('"' '"' Action42))> */ + /* 22 col <- <(( Action41) / ('\'' '\'' Action42) / ('"' '"' Action43))> */ func() bool { position247, tokenIndex247 := position, tokenIndex { @@ -2458,24 +2467,45 @@ func (p *PQL) Init() { goto l249 l250: position, tokenIndex = position249, tokenIndex249 - if buffer[position] != rune('"') { - goto l247 + if buffer[position] != rune('\'') { + goto l253 } position++ { - position253 := position - if !_rules[ruledoublequotedstring]() { - goto l247 + position254 := position + if !_rules[rulesinglequotedstring]() { + goto l253 } - add(rulePegText, position253) + add(rulePegText, position254) } - if buffer[position] != rune('"') { - goto l247 + if buffer[position] != rune('\'') { + goto l253 } position++ { add(ruleAction42, position) } + goto l249 + l253: + position, tokenIndex = position249, tokenIndex249 + if buffer[position] != rune('"') { + goto l247 + } + position++ + { + position256 := position + if !_rules[ruledoublequotedstring]() { + goto l247 + } + add(rulePegText, position256) + } + if buffer[position] != rune('"') { + goto l247 + } + position++ + { + add(ruleAction43, position) + } } l249: add(rulecol, position248) @@ -2487,99 +2517,99 @@ func (p *PQL) Init() { }, /* 23 open <- <('(' sp)> */ func() bool { - position255, tokenIndex255 := position, tokenIndex + position258, tokenIndex258 := position, tokenIndex { - position256 := position + position259 := position if buffer[position] != rune('(') { - goto l255 + goto l258 } position++ if !_rules[rulesp]() { - goto l255 + goto l258 } - add(ruleopen, position256) + add(ruleopen, position259) } return true - l255: - position, tokenIndex = position255, tokenIndex255 + l258: + position, tokenIndex = position258, tokenIndex258 return false }, /* 24 close <- <(')' sp)> */ func() bool { - position257, tokenIndex257 := position, tokenIndex + position260, tokenIndex260 := position, tokenIndex { - position258 := position + position261 := position if buffer[position] != rune(')') { - goto l257 + goto l260 } position++ if !_rules[rulesp]() { - goto l257 + goto l260 } - add(ruleclose, position258) + add(ruleclose, position261) } return true - l257: - position, tokenIndex = position257, tokenIndex257 + l260: + position, tokenIndex = position260, tokenIndex260 return false }, /* 25 sp <- <(' ' / '\t' / '\n')*> */ func() bool { { - position260 := position - l261: + position263 := position + l264: { - position262, tokenIndex262 := position, tokenIndex + position265, tokenIndex265 := position, tokenIndex { - position263, tokenIndex263 := position, tokenIndex + position266, tokenIndex266 := position, tokenIndex if buffer[position] != rune(' ') { - goto l264 + goto l267 } position++ - goto l263 - l264: - position, tokenIndex = position263, tokenIndex263 + goto l266 + l267: + position, tokenIndex = position266, tokenIndex266 if buffer[position] != rune('\t') { + goto l268 + } + position++ + goto l266 + l268: + position, tokenIndex = position266, tokenIndex266 + if buffer[position] != rune('\n') { goto l265 } position++ - goto l263 - l265: - position, tokenIndex = position263, tokenIndex263 - if buffer[position] != rune('\n') { - goto l262 - } - position++ } - l263: - goto l261 - l262: - position, tokenIndex = position262, tokenIndex262 + l266: + goto l264 + l265: + position, tokenIndex = position265, tokenIndex265 } - add(rulesp, position260) + add(rulesp, position263) } return true }, /* 26 comma <- <(sp ',' sp)> */ func() bool { - position266, tokenIndex266 := position, tokenIndex + position269, tokenIndex269 := position, tokenIndex { - position267 := position + position270 := position if !_rules[rulesp]() { - goto l266 + goto l269 } if buffer[position] != rune(',') { - goto l266 + goto l269 } position++ if !_rules[rulesp]() { - goto l266 + goto l269 } - add(rulecomma, position267) + add(rulecomma, position270) } return true - l266: - position, tokenIndex = position266, tokenIndex266 + l269: + position, tokenIndex = position269, tokenIndex269 return false }, /* 27 lbrack <- <('[' sp)> */ @@ -2590,139 +2620,139 @@ func (p *PQL) Init() { nil, /* 30 timestampbasicfmt <- <([0-9] [0-9] [0-9] [0-9] '-' ('0' / '1') [0-9] '-' [0-3] [0-9] 'T' [0-9] [0-9] ':' [0-9] [0-9])> */ func() bool { - position271, tokenIndex271 := position, tokenIndex + position274, tokenIndex274 := position, tokenIndex { - position272 := position + position275 := position if c := buffer[position]; c < rune('0') || c > rune('9') { - goto l271 + goto l274 } position++ if c := buffer[position]; c < rune('0') || c > rune('9') { - goto l271 + goto l274 } position++ if c := buffer[position]; c < rune('0') || c > rune('9') { - goto l271 + goto l274 } position++ if c := buffer[position]; c < rune('0') || c > rune('9') { - goto l271 + goto l274 } position++ if buffer[position] != rune('-') { - goto l271 + goto l274 } position++ { - position273, tokenIndex273 := position, tokenIndex + position276, tokenIndex276 := position, tokenIndex if buffer[position] != rune('0') { + goto l277 + } + position++ + goto l276 + l277: + position, tokenIndex = position276, tokenIndex276 + if buffer[position] != rune('1') { goto l274 } position++ - goto l273 - l274: - position, tokenIndex = position273, tokenIndex273 - if buffer[position] != rune('1') { - goto l271 - } - position++ } - l273: + l276: if c := buffer[position]; c < rune('0') || c > rune('9') { - goto l271 + goto l274 } position++ if buffer[position] != rune('-') { - goto l271 + goto l274 } position++ if c := buffer[position]; c < rune('0') || c > rune('3') { - goto l271 + goto l274 } position++ if c := buffer[position]; c < rune('0') || c > rune('9') { - goto l271 + goto l274 } position++ if buffer[position] != rune('T') { - goto l271 + goto l274 } position++ if c := buffer[position]; c < rune('0') || c > rune('9') { - goto l271 + goto l274 } position++ if c := buffer[position]; c < rune('0') || c > rune('9') { - goto l271 + goto l274 } position++ if buffer[position] != rune(':') { - goto l271 + goto l274 } position++ if c := buffer[position]; c < rune('0') || c > rune('9') { - goto l271 + goto l274 } position++ if c := buffer[position]; c < rune('0') || c > rune('9') { - goto l271 + goto l274 } position++ - add(ruletimestampbasicfmt, position272) + add(ruletimestampbasicfmt, position275) } return true - l271: - position, tokenIndex = position271, tokenIndex271 + l274: + position, tokenIndex = position274, tokenIndex274 return false }, /* 31 timestampfmt <- <(('"' timestampbasicfmt '"') / ('\'' timestampbasicfmt '\'') / timestampbasicfmt)> */ func() bool { - position275, tokenIndex275 := position, tokenIndex + position278, tokenIndex278 := position, tokenIndex { - position276 := position + position279 := position { - position277, tokenIndex277 := position, tokenIndex + position280, tokenIndex280 := position, tokenIndex if buffer[position] != rune('"') { - goto l278 + goto l281 } position++ if !_rules[ruletimestampbasicfmt]() { - goto l278 + goto l281 } if buffer[position] != rune('"') { + goto l281 + } + position++ + goto l280 + l281: + position, tokenIndex = position280, tokenIndex280 + if buffer[position] != rune('\'') { + goto l282 + } + position++ + if !_rules[ruletimestampbasicfmt]() { + goto l282 + } + if buffer[position] != rune('\'') { + goto l282 + } + position++ + goto l280 + l282: + position, tokenIndex = position280, tokenIndex280 + if !_rules[ruletimestampbasicfmt]() { goto l278 } - position++ - goto l277 - l278: - position, tokenIndex = position277, tokenIndex277 - if buffer[position] != rune('\'') { - goto l279 - } - position++ - if !_rules[ruletimestampbasicfmt]() { - goto l279 - } - if buffer[position] != rune('\'') { - goto l279 - } - position++ - goto l277 - l279: - position, tokenIndex = position277, tokenIndex277 - if !_rules[ruletimestampbasicfmt]() { - goto l275 - } } - l277: - add(ruletimestampfmt, position276) + l280: + add(ruletimestampfmt, position279) } return true - l275: - position, tokenIndex = position275, tokenIndex275 + l278: + position, tokenIndex = position278, tokenIndex278 return false }, - /* 32 timestamp <- <( Action43)> */ + /* 32 timestamp <- <( Action44)> */ nil, /* 34 Action0 <- <{p.startCall("Set")}> */ nil, @@ -2811,7 +2841,9 @@ func (p *PQL) Init() { nil, /* 77 Action42 <- <{p.addPosStr("_col", buffer[begin:end])}> */ nil, - /* 78 Action43 <- <{p.addPosStr("_timestamp", buffer[begin:end])}> */ + /* 78 Action43 <- <{p.addPosStr("_col", buffer[begin:end])}> */ + nil, + /* 79 Action44 <- <{p.addPosStr("_timestamp", buffer[begin:end])}> */ nil, } p.rules = _rules diff --git a/pql/pqlpeg_test.go b/pql/pqlpeg_test.go index ad40364b3..076261599 100644 --- a/pql/pqlpeg_test.go +++ b/pql/pqlpeg_test.go @@ -68,6 +68,14 @@ func TestPEGWorking(t *testing.T) { name: "Set", input: "Set(2, f=10)", ncalls: 1}, + { + name: "SetWithColKeySingleQuote", + input: `Set('foo', f=10)`, + ncalls: 1}, + { + name: "SetWithColKeyDoubleQuote", + input: `Set("foo", f=10)`, + ncalls: 1}, { name: "SetTime", input: "Set(2, f=1, 1999-12-31T00:00)", From d22507d36deecd873c1816b99e9e2a31bc161556 Mon Sep 17 00:00:00 2001 From: Travis Turner Date: Mon, 9 Jul 2018 17:05:58 -0500 Subject: [PATCH 2/2] rename gossip.NewGossipMemberSet to gossip.NewMemberSet --- gossip/gossip.go | 72 ++++++++++++++++++++++++------------------------ server/server.go | 2 +- 2 files changed, 37 insertions(+), 37 deletions(-) diff --git a/gossip/gossip.go b/gossip/gossip.go index 7650cff86..2b983376f 100644 --- a/gossip/gossip.go +++ b/gossip/gossip.go @@ -33,28 +33,28 @@ import ( ) // Ensure GossipMemberSet implements interfaces. -var _ memberlist.Delegate = &gossipMemberSet{} +var _ memberlist.Delegate = &memberSet{} -// gossipMemberSet represents a gossip implementation of MemberSet using memberlist. -type gossipMemberSet struct { +// memberSet represents a gossip implementation of MemberSet using memberlist. +type memberSet struct { mu sync.RWMutex memberlist *memberlist.Memberlist broadcasts *memberlist.TransmitLimitedQueue papi *pilosa.API - config *gossipConfig + config *config Logger pilosa.Logger logger *log.Logger transport *Transport - gossipEventReceiver *gossipEventReceiver + eventReceiver *eventReceiver } // Open implements the MemberSet interface to start network activity. -func (g *gossipMemberSet) Open() (err error) { +func (g *memberSet) Open() (err error) { g.mu.Lock() g.memberlist, err = memberlist.Create(g.config.memberlistConfig) g.mu.Unlock() @@ -94,7 +94,7 @@ func (g *gossipMemberSet) Open() (err error) { } // joinWithRetry wraps the standard memberlist Join function in a retry. -func (g *gossipMemberSet) joinWithRetry(hosts []string) error { +func (g *memberSet) joinWithRetry(hosts []string) error { err := retry(60, 2*time.Second, func() error { _, err := g.memberlist.Join(hosts) return err @@ -120,34 +120,34 @@ func retry(attempts int, sleep time.Duration, fn func() error) (err error) { //////////////////////////////////////////////////////////////// -type gossipConfig struct { +type config struct { gossipSeeds []string memberlistConfig *memberlist.Config } -// gossipMemberSetOption describes a functional option for GossipMemberSet. -type gossipMemberSetOption func(*gossipMemberSet) error +// memberSetOption describes a functional option for GossipMemberSet. +type memberSetOption func(*memberSet) error -// WithTransport is a functional option for providing a transport to NewGossipMemberSet. -func WithTransport(transport *Transport) gossipMemberSetOption { - return func(g *gossipMemberSet) error { +// WithTransport is a functional option for providing a transport to NewMemberSet. +func WithTransport(transport *Transport) memberSetOption { + return func(g *memberSet) error { g.transport = transport return nil } } -// WithLogger is a functional option for providing a logger to NewGossipMemberSet. -func WithLogger(logger *log.Logger) gossipMemberSetOption { - return func(g *gossipMemberSet) error { +// WithLogger is a functional option for providing a logger to NewMemberSet. +func WithLogger(logger *log.Logger) memberSetOption { + return func(g *memberSet) error { g.logger = logger return nil } } -// NewGossipMemberSet returns a new instance of GossipMemberSet based on options. -func NewGossipMemberSet(cfg Config, api *pilosa.API, options ...gossipMemberSetOption) (*gossipMemberSet, error) { +// NewMemberSet returns a new instance of GossipMemberSet based on options. +func NewMemberSet(cfg Config, api *pilosa.API, options ...memberSetOption) (*memberSet, error) { host := api.Node().URI.Host - g := &gossipMemberSet{ + g := &memberSet{ papi: api, Logger: pilosa.NopLogger, } @@ -158,8 +158,8 @@ func NewGossipMemberSet(cfg Config, api *pilosa.API, options ...gossipMemberSetO return nil, errors.Wrap(err, "executing option") } } - ger := newGossipEventReceiver(g.logger, api) - g.gossipEventReceiver = ger + ger := newEventReceiver(g.logger, api) + g.eventReceiver = ger if g.transport == nil { port, err := strconv.Atoi(cfg.Port) @@ -210,7 +210,7 @@ func NewGossipMemberSet(cfg Config, api *pilosa.API, options ...gossipMemberSetO conf.Events = ger conf.Logger = g.logger - g.config = &gossipConfig{ + g.config = &config{ memberlistConfig: conf, gossipSeeds: cfg.Seeds, } @@ -219,7 +219,7 @@ func NewGossipMemberSet(cfg Config, api *pilosa.API, options ...gossipMemberSetO } // NodeMeta implementation of the memberlist.Delegate interface. -func (g *gossipMemberSet) NodeMeta(limit int) []byte { +func (g *memberSet) NodeMeta(limit int) []byte { buf, err := g.papi.Serializer.Marshal(g.papi.Node()) if err != nil { g.Logger.Printf("marshal message error: %s", err) @@ -230,7 +230,7 @@ func (g *gossipMemberSet) NodeMeta(limit int) []byte { // NotifyMsg implementation of the memberlist.Delegate interface // called when a user-data message is received. -func (g *gossipMemberSet) NotifyMsg(b []byte) { +func (g *memberSet) NotifyMsg(b []byte) { err := g.papi.ClusterMessage(context.Background(), bytes.NewBuffer(b)) if err != nil { g.Logger.Printf("cluster message error: %s", err) @@ -239,13 +239,13 @@ func (g *gossipMemberSet) NotifyMsg(b []byte) { // GetBroadcasts implementation of the memberlist.Delegate interface // called when user data messages can be broadcast. -func (g *gossipMemberSet) GetBroadcasts(overhead, limit int) [][]byte { +func (g *memberSet) 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 *gossipMemberSet) LocalState(join bool) []byte { +func (g *memberSet) LocalState(join bool) []byte { m := &pilosa.NodeStatus{ Node: g.papi.Node(), MaxShards: g.papi.MaxShards(context.Background()), @@ -263,28 +263,28 @@ func (g *gossipMemberSet) LocalState(join bool) []byte { // MergeRemoteState implementation of the memberlist.Delegate interface // receive and process the remote side's LocalState. -func (g *gossipMemberSet) MergeRemoteState(buf []byte, join bool) { +func (g *memberSet) MergeRemoteState(buf []byte, join bool) { err := g.papi.ClusterMessage(context.Background(), bytes.NewBuffer(buf)) if err != nil { g.Logger.Printf("merge state error: %s", err) } } -// gossipEventReceiver is used to enable an application to receive +// eventReceiver is used to enable an application to receive // events about joins and leaves over a channel. // // Care must be taken that events are processed in a timely manner from // the channel, since this delegate will block until an event can be sent. -type gossipEventReceiver struct { +type eventReceiver struct { ch chan memberlist.NodeEvent papi *pilosa.API logger *log.Logger } -// newGossipEventReceiver returns a new instance of GossipEventReceiver. -func newGossipEventReceiver(logger *log.Logger, papi *pilosa.API) *gossipEventReceiver { - ger := &gossipEventReceiver{ +// newEventReceiver returns a new instance of GossipEventReceiver. +func newEventReceiver(logger *log.Logger, papi *pilosa.API) *eventReceiver { + ger := &eventReceiver{ ch: make(chan memberlist.NodeEvent, 1), logger: logger, papi: papi, @@ -293,19 +293,19 @@ func newGossipEventReceiver(logger *log.Logger, papi *pilosa.API) *gossipEventRe return ger } -func (g *gossipEventReceiver) NotifyJoin(n *memberlist.Node) { +func (g *eventReceiver) NotifyJoin(n *memberlist.Node) { g.ch <- memberlist.NodeEvent{memberlist.NodeJoin, n} } -func (g *gossipEventReceiver) NotifyLeave(n *memberlist.Node) { +func (g *eventReceiver) NotifyLeave(n *memberlist.Node) { g.ch <- memberlist.NodeEvent{memberlist.NodeLeave, n} } -func (g *gossipEventReceiver) NotifyUpdate(n *memberlist.Node) { +func (g *eventReceiver) NotifyUpdate(n *memberlist.Node) { g.ch <- memberlist.NodeEvent{memberlist.NodeUpdate, n} } -func (g *gossipEventReceiver) listen() { +func (g *eventReceiver) listen() { var nodeEventType pilosa.NodeEventType for { e := <-g.ch diff --git a/server/server.go b/server/server.go index d71139096..e83cc1c1c 100644 --- a/server/server.go +++ b/server/server.go @@ -317,7 +317,7 @@ func (m *Command) setupNetworking() error { return errors.Wrap(err, "getting transport") } - gossipMemberSet, err := gossip.NewGossipMemberSet( + gossipMemberSet, err := gossip.NewMemberSet( m.Config.Gossip, m.API, gossip.WithLogger(m.logger.Logger()),