From adf3e528f56ed6da5a6b29b644b5c8c0d82adb81 Mon Sep 17 00:00:00 2001 From: reesporte Date: Thu, 21 Oct 2021 13:09:24 -0500 Subject: [PATCH 01/10] Change time estimation to use avg time per message MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit In [SUP-75](https://molecula.atlassian.net/browse/SUP-75?atlOrigin=eyJpIjoiYmU5MzdkMmUyZTAyNGQ2Y2IzMDMzYTgzMDU2Y2ZhNmMiLCJwIjoiaiJ9) Allen pointed out that the time estimation is really good for the first couple lines of output, but gets exponentially worse as execution continues. After looking into it, it looks like we’re currently using a heuristic based on the amount of messages processed in the previous second(ish) which is what results in that sort of exponential drop off. To remedy this, I adjusted the time estimation calculation to use the average time per message up to the point of calculating the new estimate to ideally improve estimates over time, with the trade-off of a potentially less accurate estimate to begin with. --- server.go | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/server.go b/server.go index fdca4dc39..ce6babe25 100644 --- a/server.go +++ b/server.go @@ -722,8 +722,10 @@ func (s *Server) Open() error { if now := time.Now(); now.Sub(prevMsg) > time.Second { progressRatio := float64(i+1) / float64(len(toSend)) - remainingRatio := 1 - progressRatio - timeRemaining := time.Duration(float64(now.Sub(prevMsg)) * (remainingRatio / progressRatio)) + numSentMessages := len(toSend) - (i + 1) + messagesLeft := len(toSend) - numSentMessages + avgTimePerMessage := float64(now.Sub(start)) / float64(numSentMessages) + timeRemaining := time.Duration(avgTimePerMessage * float64(messagesLeft)) s.logger.Printf("synced %d/%d messages (%.2f%% complete; %s remaining)", i+1, len(toSend), 100*progressRatio, timeRemaining) prevMsg = now } From 27515a3f98d5518f5d85188fce7995fdc9986617 Mon Sep 17 00:00:00 2001 From: reesporte Date: Thu, 21 Oct 2021 14:24:27 -0500 Subject: [PATCH 02/10] number of sent messages is just i silly --- server.go | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/server.go b/server.go index ce6babe25..8ee1ba5bc 100644 --- a/server.go +++ b/server.go @@ -722,9 +722,8 @@ func (s *Server) Open() error { if now := time.Now(); now.Sub(prevMsg) > time.Second { progressRatio := float64(i+1) / float64(len(toSend)) - numSentMessages := len(toSend) - (i + 1) - messagesLeft := len(toSend) - numSentMessages - avgTimePerMessage := float64(now.Sub(start)) / float64(numSentMessages) + messagesLeft := len(toSend) - i + avgTimePerMessage := float64(now.Sub(start)) / float64(i) timeRemaining := time.Duration(avgTimePerMessage * float64(messagesLeft)) s.logger.Printf("synced %d/%d messages (%.2f%% complete; %s remaining)", i+1, len(toSend), 100*progressRatio, timeRemaining) prevMsg = now From 7102de9f605c7e2900ef8589634edebee1d44dd5 Mon Sep 17 00:00:00 2001 From: reesporte Date: Thu, 21 Oct 2021 16:19:57 -0500 Subject: [PATCH 03/10] off by one error fixed --- server.go | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/server.go b/server.go index 8ee1ba5bc..44e9da6f2 100644 --- a/server.go +++ b/server.go @@ -722,8 +722,8 @@ func (s *Server) Open() error { if now := time.Now(); now.Sub(prevMsg) > time.Second { progressRatio := float64(i+1) / float64(len(toSend)) - messagesLeft := len(toSend) - i - avgTimePerMessage := float64(now.Sub(start)) / float64(i) + messagesLeft := len(toSend) - (i + 1) + avgTimePerMessage := float64(now.Sub(start)) / float64(i+1) timeRemaining := time.Duration(avgTimePerMessage * float64(messagesLeft)) s.logger.Printf("synced %d/%d messages (%.2f%% complete; %s remaining)", i+1, len(toSend), 100*progressRatio, timeRemaining) prevMsg = now From 45e36600f3b299edea7544c295975991456f3bac Mon Sep 17 00:00:00 2001 From: reesporte Date: Fri, 22 Oct 2021 10:49:07 -0500 Subject: [PATCH 04/10] don't include .*.swp --- .gitignore | 1 + 1 file changed, 1 insertion(+) diff --git a/.gitignore b/.gitignore index d53a2e53c..21dbedb4d 100644 --- a/.gitignore +++ b/.gitignore @@ -11,3 +11,4 @@ release-pilosa-fsck.*.*.tar.gz pilosa *.dot .idea/ +.*.swp From 0f108a612dc5ff23b11d6248b543795fa0e74239 Mon Sep 17 00:00:00 2001 From: reesporte Date: Fri, 22 Oct 2021 11:25:03 -0500 Subject: [PATCH 05/10] refactor and add unit tests --- server.go | 8 +++---- util.go | 9 ++++++++ util_test.go | 59 ++++++++++++++++++++++++++++++++++++++++++++++++++++ 3 files changed, 71 insertions(+), 5 deletions(-) create mode 100644 util_test.go diff --git a/server.go b/server.go index 44e9da6f2..f6fd1ef5b 100644 --- a/server.go +++ b/server.go @@ -703,6 +703,7 @@ func (s *Server) Open() error { start := time.Now() prevMsg := start + numMsgs := uint(len(toSend)) s.logger.Printf("start initial cluster state sync") for i := range toSend { for { @@ -721,11 +722,8 @@ func (s *Server) Open() error { } if now := time.Now(); now.Sub(prevMsg) > time.Second { - progressRatio := float64(i+1) / float64(len(toSend)) - messagesLeft := len(toSend) - (i + 1) - avgTimePerMessage := float64(now.Sub(start)) / float64(i+1) - timeRemaining := time.Duration(avgTimePerMessage * float64(messagesLeft)) - s.logger.Printf("synced %d/%d messages (%.2f%% complete; %s remaining)", i+1, len(toSend), 100*progressRatio, timeRemaining) + pctDone := (float64(i+1) / float64(numMsgs)) * 100 + s.logger.Printf("synced %d/%d messages (%.2f%% complete; %s remaining)", i+1, numMsgs, pctDone, EstTimeLeft(start, now, uint(i), numMsgs)) prevMsg = now } } diff --git a/util.go b/util.go index 3bbeff8e3..93f5e9878 100644 --- a/util.go +++ b/util.go @@ -25,6 +25,7 @@ import ( "sort" "strings" "syscall" + "time" "unsafe" "github.com/molecula/featurebase/v2/roaring" @@ -343,3 +344,11 @@ func roaringFragmentHasData(path string, index, field, view string, shard uint64 return } + +// EstTimeLeft returns the estimated remaining time to iterate through some items +// given a start time, the current time, the iteration, and the number of items +func EstTimeLeft(start time.Time, now time.Time, i uint, total uint) time.Duration { + msgsLeft := total - (i + 1) + avgMsgTime := float64(now.Sub(start)) / float64(i+1) + return time.Duration(avgMsgTime * float64(msgsLeft)) +} diff --git a/util_test.go b/util_test.go new file mode 100644 index 000000000..824035744 --- /dev/null +++ b/util_test.go @@ -0,0 +1,59 @@ +package pilosa + +// util_test.go has unit tests for utility functions from util.go + +import ( + "testing" + "time" +) + +func TestEstTimeLeft(t *testing.T) { + cases := []struct { + start time.Time + now time.Time + i uint + total uint + }{ + { + time.Date(1969, time.June, 9, 4, 20, 0, 0, time.UTC), + time.Date(1969, time.June, 9, 4, 21, 0, 0, time.UTC), + 10, + 20, + }, + { + time.Date(1969, time.June, 9, 4, 20, 0, 0, time.UTC), + time.Date(1969, time.June, 9, 4, 20, 4, 0, time.UTC), + 10, + 20, + }, + { + time.Date(1969, time.June, 9, 4, 20, 0, 0, time.UTC), + time.Date(1969, time.June, 9, 4, 21, 1, 5, time.UTC), + 10, + 20, + }, + { + time.Date(1969, time.June, 9, 4, 20, 0, 0, time.UTC), + time.Date(1969, time.June, 9, 4, 21, 0, 0, time.UTC), + 1, + 20, + }, + { + time.Date(1969, time.June, 9, 4, 20, 0, 0, time.UTC), + time.Date(1969, time.June, 9, 4, 21, 0, 0, time.UTC), + 10, + 50, + }, + } + + for _, c := range cases { + // we expect that it will be the avg time per message times + // the number of remaining messages + expected := time.Duration((float64(c.now.Sub(c.start)) / float64(c.i+1)) * float64(c.total-(c.i+1))) + + timeLeft := EstTimeLeft(c.start, c.now, c.i, c.total) + if timeLeft != expected { + t.Errorf("Time left was incorrect, expected: %d, but got: %d", expected, timeLeft) + } + } +} From 681ed9923d16804f932963b8cdbf0bf9ec810638 Mon Sep 17 00:00:00 2001 From: reesporte Date: Fri, 22 Oct 2021 11:29:09 -0500 Subject: [PATCH 06/10] add license header --- util_test.go | 14 ++++++++++++++ 1 file changed, 14 insertions(+) diff --git a/util_test.go b/util_test.go index 824035744..9aa336f31 100644 --- a/util_test.go +++ b/util_test.go @@ -1,3 +1,17 @@ +// Copyright 2021 Pilosa Corp. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + package pilosa // util_test.go has unit tests for utility functions from util.go From bd0d68b2fd4f3c7d47fb3d270485abba41f1567f Mon Sep 17 00:00:00 2001 From: reesporte Date: Fri, 22 Oct 2021 12:54:46 -0500 Subject: [PATCH 07/10] rename function, return pctDone --- server.go | 4 ++-- util.go | 10 ++++++---- util_test.go | 8 ++++++-- 3 files changed, 14 insertions(+), 8 deletions(-) diff --git a/server.go b/server.go index f6fd1ef5b..8912d38bb 100644 --- a/server.go +++ b/server.go @@ -722,8 +722,8 @@ func (s *Server) Open() error { } if now := time.Now(); now.Sub(prevMsg) > time.Second { - pctDone := (float64(i+1) / float64(numMsgs)) * 100 - s.logger.Printf("synced %d/%d messages (%.2f%% complete; %s remaining)", i+1, numMsgs, pctDone, EstTimeLeft(start, now, uint(i), numMsgs)) + estimate, pctDone := GetLoopProgress(start, now, uint(i), numMsgs) + s.logger.Printf("synced %d/%d messages (%.2f%% complete; %s remaining)", i+1, numMsgs, pctDone, estimate) prevMsg = now } } diff --git a/util.go b/util.go index 93f5e9878..d165f634e 100644 --- a/util.go +++ b/util.go @@ -345,10 +345,12 @@ func roaringFragmentHasData(path string, index, field, view string, shard uint64 return } -// EstTimeLeft returns the estimated remaining time to iterate through some items -// given a start time, the current time, the iteration, and the number of items -func EstTimeLeft(start time.Time, now time.Time, i uint, total uint) time.Duration { +// GetLoopProgress returns the estimated remaining time to iterate through some items +// as well as the loop completion percentage with the following parameters: +// the start time, the current time, the iteration, and the number of items +func GetLoopProgress(start time.Time, now time.Time, i uint, total uint) (time.Duration, float64) { msgsLeft := total - (i + 1) avgMsgTime := float64(now.Sub(start)) / float64(i+1) - return time.Duration(avgMsgTime * float64(msgsLeft)) + pctDone := (float64(i+1) / float64(total)) * 100 + return time.Duration(avgMsgTime * float64(msgsLeft)), pctDone } diff --git a/util_test.go b/util_test.go index 9aa336f31..ea27f1c5b 100644 --- a/util_test.go +++ b/util_test.go @@ -21,7 +21,7 @@ import ( "time" ) -func TestEstTimeLeft(t *testing.T) { +func TestGetLoopProgress(t *testing.T) { cases := []struct { start time.Time now time.Time @@ -64,10 +64,14 @@ func TestEstTimeLeft(t *testing.T) { // we expect that it will be the avg time per message times // the number of remaining messages expected := time.Duration((float64(c.now.Sub(c.start)) / float64(c.i+1)) * float64(c.total-(c.i+1))) + expectedPct := 100 * (float64(c.i+1) / float64(c.total)) - timeLeft := EstTimeLeft(c.start, c.now, c.i, c.total) + timeLeft, pctDone := GetLoopProgress(c.start, c.now, c.i, c.total) if timeLeft != expected { t.Errorf("Time left was incorrect, expected: %d, but got: %d", expected, timeLeft) } + if pctDone != expectedPct { + t.Errorf("Percentage done was incorrect, expected: %f, but got: %f", expectedPct, pctDone) + } } } From 5dd4b0e048a6e2f6ab4ac4673639930f857d4a25 Mon Sep 17 00:00:00 2001 From: reesporte Date: Fri, 22 Oct 2021 14:49:05 -0500 Subject: [PATCH 08/10] rename vars to more sensible names --- util.go | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/util.go b/util.go index d165f634e..ebbff6e75 100644 --- a/util.go +++ b/util.go @@ -348,9 +348,9 @@ func roaringFragmentHasData(path string, index, field, view string, shard uint64 // GetLoopProgress returns the estimated remaining time to iterate through some items // as well as the loop completion percentage with the following parameters: // the start time, the current time, the iteration, and the number of items -func GetLoopProgress(start time.Time, now time.Time, i uint, total uint) (time.Duration, float64) { - msgsLeft := total - (i + 1) - avgMsgTime := float64(now.Sub(start)) / float64(i+1) - pctDone := (float64(i+1) / float64(total)) * 100 - return time.Duration(avgMsgTime * float64(msgsLeft)), pctDone +func GetLoopProgress(start time.Time, now time.Time, iteration uint, total uint) (remaining time.Duration, pctDone float64) { + itemsLeft := total - (iteration + 1) + avgItemTime := float64(now.Sub(start)) / float64(iteration+1) + pctDone = (float64(iteration+1) / float64(total)) * 100 + return time.Duration(avgItemTime * float64(itemsLeft)), pctDone } From d934d117da74875477b0400460c8103aefe90548 Mon Sep 17 00:00:00 2001 From: reesporte Date: Fri, 22 Oct 2021 14:52:57 -0500 Subject: [PATCH 09/10] add test case names --- util_test.go | 31 ++++++++++++++++++++----------- 1 file changed, 20 insertions(+), 11 deletions(-) diff --git a/util_test.go b/util_test.go index ea27f1c5b..3416a95a8 100644 --- a/util_test.go +++ b/util_test.go @@ -22,37 +22,44 @@ import ( ) func TestGetLoopProgress(t *testing.T) { + // TODO: try to find more sneaky cases cases := []struct { + name string start time.Time now time.Time i uint total uint }{ { + "one minute, half done", time.Date(1969, time.June, 9, 4, 20, 0, 0, time.UTC), time.Date(1969, time.June, 9, 4, 21, 0, 0, time.UTC), 10, 20, }, { + "four seconds, half done", time.Date(1969, time.June, 9, 4, 20, 0, 0, time.UTC), time.Date(1969, time.June, 9, 4, 20, 4, 0, time.UTC), 10, 20, }, { + "one minute, one second, 5 μs, half done", time.Date(1969, time.June, 9, 4, 20, 0, 0, time.UTC), time.Date(1969, time.June, 9, 4, 21, 1, 5, time.UTC), 10, 20, }, { + "one minute, one done", time.Date(1969, time.June, 9, 4, 20, 0, 0, time.UTC), time.Date(1969, time.June, 9, 4, 21, 0, 0, time.UTC), 1, 20, }, { + "one minute, 1/5 done", time.Date(1969, time.June, 9, 4, 20, 0, 0, time.UTC), time.Date(1969, time.June, 9, 4, 21, 0, 0, time.UTC), 10, @@ -61,17 +68,19 @@ func TestGetLoopProgress(t *testing.T) { } for _, c := range cases { - // we expect that it will be the avg time per message times - // the number of remaining messages - expected := time.Duration((float64(c.now.Sub(c.start)) / float64(c.i+1)) * float64(c.total-(c.i+1))) - expectedPct := 100 * (float64(c.i+1) / float64(c.total)) + t.Run(c.name, func(t *testing.T) { + // we expect that it will be the avg time per message times + // the number of remaining messages + expected := time.Duration((float64(c.now.Sub(c.start)) / float64(c.i+1)) * float64(c.total-(c.i+1))) + expectedPct := 100 * (float64(c.i+1) / float64(c.total)) - timeLeft, pctDone := GetLoopProgress(c.start, c.now, c.i, c.total) - if timeLeft != expected { - t.Errorf("Time left was incorrect, expected: %d, but got: %d", expected, timeLeft) - } - if pctDone != expectedPct { - t.Errorf("Percentage done was incorrect, expected: %f, but got: %f", expectedPct, pctDone) - } + timeLeft, pctDone := GetLoopProgress(c.start, c.now, c.i, c.total) + if timeLeft != expected { + t.Errorf("Time left was incorrect, expected: %d, but got: %d", expected, timeLeft) + } + if pctDone != expectedPct { + t.Errorf("Percentage done was incorrect, expected: %f, but got: %f", expectedPct, pctDone) + } + }) } } From d51c6b950f1cf05253ae481b86b6ea9cfc920c20 Mon Sep 17 00:00:00 2001 From: reesporte Date: Mon, 25 Oct 2021 15:24:55 -0500 Subject: [PATCH 10/10] meaningless commit to kick off sonarcloud with new rules --- util_test.go | 1 + 1 file changed, 1 insertion(+) diff --git a/util_test.go b/util_test.go index 3416a95a8..c659053c5 100644 --- a/util_test.go +++ b/util_test.go @@ -15,6 +15,7 @@ package pilosa // util_test.go has unit tests for utility functions from util.go +// import ( "testing"