From 08d4511e6c5da0ba399259579601cd1ff0c8237e Mon Sep 17 00:00:00 2001 From: Matt Jaffee Date: Fri, 5 Mar 2021 12:51:39 -0600 Subject: [PATCH 1/2] fix issue where restarting a cluster logged "invalid pilosa.Message" we were sending pilosa.Message objects from a spool, but actually passing a pointer to them rather than the Message itself. I'm concerned this wasn't caught in any test, and also curious if that needed to be a pointer for some reason or if it's a typo. Definitely need to write a test still. --- server.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/server.go b/server.go index bd08550f0..fcd70ef13 100644 --- a/server.go +++ b/server.go @@ -679,7 +679,7 @@ func (s *Server) Open() error { s.logger.Printf("start initial cluster state sync") for i := range toSend { for { - err := s.holder.broadcaster.SendSync(&toSend[i]) + err := s.holder.broadcaster.SendSync(toSend[i]) if err != nil { s.logger.Printf("failed to broadcast startup cluster message (trying again in a bit): %v", err) timer.Reset(time.Second) From f97ac4928c568143a783b25d4d9cf19ba59d2a6b Mon Sep 17 00:00:00 2001 From: Matt Jaffee Date: Fri, 5 Mar 2021 14:16:52 -0600 Subject: [PATCH 2/2] add test for pilosa.Message issue this doesn't quite work as-is, but I verified that it reproduced/fixed the issue by adding a panic where the problem log statement is. There's a follow up ticket to fix the test... it's just a bit open ended as to the best way to do that. --- executor_test.go | 52 ++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 52 insertions(+) diff --git a/executor_test.go b/executor_test.go index 7e79baa14..5b64bf990 100644 --- a/executor_test.go +++ b/executor_test.go @@ -3512,6 +3512,58 @@ func TestExecutor_ExecuteOptions(t *testing.T) { }) } +func TestReopenCluster(t *testing.T) { + commandOpts := make([][]server.CommandOption, 3) + configs := make([]*server.Config, 3) + for i := range configs { + conf := server.NewConfig() + configs[i] = conf + conf.Cluster.ReplicaN = 2 + commandOpts[i] = append(commandOpts[i], server.OptCommandConfig(conf)) + } + c := test.MustRunCluster(t, 3, commandOpts...) + defer c.Close() + c.CreateField(t, "users", pilosa.IndexOptions{Keys: true, TrackExistence: true}, "likenums") + c.ImportIDKey(t, "users", "likenums", []test.KeyID{ + {ID: 1, Key: "userA"}, + {ID: 2, Key: "userB"}, + {ID: 3, Key: "userC"}, + {ID: 4, Key: "userD"}, + {ID: 5, Key: "userE"}, + {ID: 6, Key: "userF"}, + {ID: 7, Key: "userA"}, + {ID: 7, Key: "userB"}, + {ID: 7, Key: "userC"}, + {ID: 7, Key: "userD"}, + // we intentionally leave user E out because then there is no + // data for userE's shard for this field, which triggered a + // "fragment not found" problem + {ID: 7, Key: "userF"}, + }) + + node0 := c.GetNode(0) + if err := node0.Reopen(); err != nil { + t.Fatal(err) + } + if err := node0.AwaitState(disco.ClusterStateNormal, 10*time.Second); err != nil { + t.Fatalf("restarting cluster: %v", err) + } + + // TODO this test was supposed to reproduce an issue where spooled messages got sent improperly in Server.Open in this area: + // + // for i := range toSend { + // for { + // err := s.holder.broadcaster.SendSync(toSend[i]) + // + // It was previously sending &toSend[i] which wasn't a valid + // pilosa.Message. Unfortunately because that's all happening in a + // goroutine, the Reopen continues happily and things appear to + // work whether the bug is present or not. I'd like to pass in a + // logger and detect when "completed initial cluster state sync" + // is logged, but it's more involved than I have time for at the + // moment, so I'm creating a follow up ticket. (CORE-318) +} + // Ensure an existence field is maintained. func TestExecutor_Execute_Existence(t *testing.T) { t.Run("Row", func(t *testing.T) {