mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-09-15 08:41:02 +00:00
* use mainline etcd-io dependencies, not forks
this commit does lots of things around clustering with goal of increasing stability.
- upgrades from molecula/etcd to go.etcd.io/etcd@v3.5.4
- upgrades from seebs/bbolt to go.etcd.io/bbolt@v1.3.6
- update tests to use unix sockets for etcd cluster communication
- this is what etcd uses for a lot of internal testing, so if their devs think
it's a valid test, we can probably accept that
- cleanup etcd node-watcher shutdown process
Co-authored-by: tgruben <tgruben@gmail.com>
* moved random query to another repo
it had weird dependency issues with upgrading to mainline etcd bc of the vegeta dep
so we removed it bc no one really uses it anyway
we got this error message:
```
github.com/molecula/featurebase/v3/cmd/random-query imports
github.com/tsenart/vegeta/v12/lib tested by
github.com/tsenart/vegeta/v12/lib.test imports
github.com/streadway/quantile tested by
github.com/streadway/quantile.test imports
.: "." is relative, but relative import paths are not supported in module mode
```
* add cleanup to EtcdUnixSocket test util
Co-authored-by: reesporte <reesedporter@gmail.com>
125 lines
3.7 KiB
Go
125 lines
3.7 KiB
Go
// Copyright 2021 Molecula Corp. All rights reserved.
|
|
package pilosa
|
|
|
|
// util.go: a place for generic, reusable utilities.
|
|
import (
|
|
"fmt"
|
|
"os"
|
|
"os/user"
|
|
"path/filepath"
|
|
"reflect"
|
|
"strings"
|
|
"sync"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/shirou/gopsutil/v3/mem"
|
|
)
|
|
|
|
var clientPort = os.Getpid()
|
|
var muClientPort = sync.Mutex{}
|
|
|
|
// LeftShifted16MaxContainerKey is 0xffffffffffff0000. It is similar
|
|
// to the roaring.maxContainerKey 0x0000ffffffffffff, but
|
|
// shifted 16 bits to the left so its domain is the full [0, 2^64) bit space.
|
|
// It is used to match the semantics of the roaring.OffsetRange() API.
|
|
// This is the maximum endx value for Tx.OffsetRange(), because the lowbits,
|
|
// as in the roaring.OffsetRange(), are not allowed to be set.
|
|
// It is used in Tx.RoaringBitamp() to obtain the full contents of a fragment
|
|
// from a call from tx.OffsetRange() by requesting [0, LeftShifted16MaxContainerKey)
|
|
// with an offset of 0.
|
|
const LeftShifted16MaxContainerKey = uint64(0xffffffffffff0000) // or math.MaxUint64 - (1<<16 - 1), or 18446744073709486080
|
|
|
|
// NilInside checks if the provided iface is nil or
|
|
// contains a nil pointer, slice, array, map, or channel.
|
|
func NilInside(iface interface{}) bool {
|
|
if iface == nil {
|
|
return true
|
|
}
|
|
switch reflect.TypeOf(iface).Kind() {
|
|
case reflect.Ptr, reflect.Slice, reflect.Array, reflect.Map, reflect.Chan:
|
|
return reflect.ValueOf(iface).IsNil()
|
|
}
|
|
return false
|
|
}
|
|
|
|
//////////////////////////////////
|
|
// helper utility functions
|
|
|
|
func highbits(v uint64) uint64 { return v >> 16 }
|
|
func lowbits(v uint64) uint16 { return uint16(v & 0xFFFF) }
|
|
|
|
// 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, 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
|
|
}
|
|
|
|
// FormatTimestampNano returns the string representation of a timestamp given:
|
|
// an epoch value, base, and time unit
|
|
func FormatTimestampNano(value, base int64, timeUnit string) string {
|
|
return time.Unix(0, (value+base)*TimeUnitNanos(timeUnit)).UTC().Format(time.RFC3339Nano)
|
|
}
|
|
|
|
type MemoryUsage struct {
|
|
Capacity uint64 `json:"capacity"`
|
|
TotalUse uint64 `json:"totalUsed"`
|
|
}
|
|
|
|
// GetMemoryUsage gets the memory usage
|
|
func GetMemoryUsage() (MemoryUsage, error) {
|
|
usage, err := mem.VirtualMemory()
|
|
if usage == nil || err != nil {
|
|
return MemoryUsage{}, fmt.Errorf("reading virtual memory: %v", err)
|
|
}
|
|
return MemoryUsage{Capacity: usage.Total, TotalUse: usage.Used}, nil
|
|
}
|
|
|
|
type DiskUsage struct {
|
|
Usage int64 `json:"usage"`
|
|
}
|
|
|
|
// GetDiskUsage gets the disk usage of the path
|
|
func GetDiskUsage(path string) (DiskUsage, error) {
|
|
usr, _ := user.Current()
|
|
dir := usr.HomeDir
|
|
if path == "~" {
|
|
path = dir
|
|
} else if strings.HasPrefix(path, "~/") {
|
|
path = filepath.Join(dir, path[2:])
|
|
}
|
|
|
|
var size int64
|
|
err := filepath.Walk(path, func(_ string, info os.FileInfo, err error) error {
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if !info.IsDir() {
|
|
size += info.Size()
|
|
}
|
|
return err
|
|
})
|
|
return DiskUsage{size}, err
|
|
}
|
|
|
|
// EtcdUnixSocket returns a url for use in test etcd clusters.
|
|
func EtcdUnixSocket(tb testing.TB) string {
|
|
muClientPort.Lock()
|
|
defer func() {
|
|
clientPort++
|
|
muClientPort.Unlock()
|
|
}()
|
|
addr := fmt.Sprintf("fake:%d", clientPort)
|
|
tb.Cleanup(func() {
|
|
err := os.Remove(addr)
|
|
if err != nil {
|
|
tb.Logf("could not remove '%s', %v", addr, err)
|
|
}
|
|
})
|
|
return fmt.Sprintf("unix://%s", addr)
|
|
}
|