Merged with master

This commit is contained in:
Yuce Tekol 2018-11-21 18:16:43 +03:00
commit ac91635628
No known key found for this signature in database
GPG key ID: CB59E46D2FB90573
48 changed files with 3668 additions and 1074 deletions

View file

@ -44,6 +44,12 @@ jobs:
<<: *base-test
environment:
GOARCH: 386
cluster-tests:
<<: *defaults
steps:
- *fast-checkout
- setup_remote_docker
- run: make clustertests-build
prerelease:
<<: *base-test
steps:
@ -99,6 +105,9 @@ workflows:
- test-golang-1.10-386:
requires:
- build
- cluster-tests:
requires:
- build
- prerelease:
requires:
- linter

View file

@ -1,52 +0,0 @@
language: go
go:
- "1.10" # Use string, as 1.10==1.1 if interpreted as float.
- master
env:
global: # AWS_ACCESS_KEY_ID and AWS_SECRET_ACCESS_KEY
- secure: "VnBFmFfBOrrf7ONLN9WpAFCcV8SEt5G5VPnnHv97TP7PlJG8LWR6k6O+vRJOvf8V4vDMfKCTDonwWLgbssVf3yygo3C8ZoftY2phehEkWGffCgsd9ML/YBNbGq4LYLSE5HKvBqrZjQaOrVby71BAsP8W7RhC6hqzFQ00M/z8dZVfwaQQFwew2eEcSxLEaaDFS8Wgc3/UuwxDRPBq6u3cCN5RxfB+q70HvGVq4TT+0dqS4eCvz688+Z0GIGYx9olNjh0F2Kc8R2Po0lnUNa0GiHrZ21zeQ1DxIK04QABrWWmjL4h+bx3VHNKPFR4GYSKDf+pj1kfaqbfrAg6rMAJdGejgoS+QyjhgCoN4d3qRp8s+1nrxtp0TvezEdjwyxt4quGHbP5TxWUszssbGhWqf4mx6OeJ8MmdTaJjfu0f3NWJXMycqT6J73WKORk4rHeIqF9CIdxdmcpkwYj8rk0TEMTPTsd7WA8w2HIDsCz/jQnRmEgLUiNnTAofYc/uUi/Wg/T2hllkp+oBDTzxk9NTelkqx8TJ0bDmYYL9JWUi1siFHTHiVYTJgyirSfGNpe61u8OLmT0Hak/D399IfL7qgFLlMXk8q92typfO2xEduq6G+8KygeqiOMSsOY+xcDvZf5xtcEihYd21vjtrxRSqFsup/o8DIxEurQnfXBx1B+WA="
- secure: "U4fpHWDVOG4viqZsiVgUDW7OW1JW60uPOZy0q9pfbs86iHvmZq0PaScsZ+YdlYaN2GETVr7endDf6DCcZs1PWfg0F6VQfkOXcShX8HVS9O58lUZA5tyvbDVql9DQs4PbnkZo+ktz+Z0YaXqq2RdtMDOUz4bgZwspLPMA14if+N6w0tqCFpB7bEtpptTGsdbIQPG1n07yvSeNmK4mvrEEs77tWmhulN5iilpOqhpIvD39bJvtCYVALuJpzLd/OjLTPV9l/fl+hJkMXSj+X5ilO1DHINAcCM648iEX2phXAIWmi0O0Rbg2cI4kV9T5ysOIw8ux+YCm9bZDGTCt+VGBW5Fg+Z5iaXXexyKYCGiHleOJ7kCj9kXxh2u8NiYVNgb19dGJV5/HgQ6pcGWjeVEqr8yY1546zMjpTX+SYGQF+XZe+uggEjeAsk53ueXa0pyZTrlrqSvR7BBtWPx47s/dTg2L19FQYv3XpGMxEXLw92RplExQKi1h7QgihRxFpjGgURHhrt7d9eiNiNqBt3ZsHjmh2AkXZHnaDjlgSnFFWaMqP3UtDBWIuO+2BMbZUJVfP+gpQGBZ4gtpUSmV2JDCHgZgX5OAnLD4usxh+ATQ4rvUXF/tf8nMqEKHlGKd8hxpYSyMX21BoqfSfY4/IA0ejVE9BITqlrvqewqkP1yxe7o="
matrix:
- GOARCH=386
- GOARCH=386 ENTERPRISE=1
- GOARCH=amd64
- GOARCH=amd64 ENTERPRISE=1
cache:
directories:
vendor
install:
- make -B install-dep vendor
script: make test
jobs:
include:
- stage: metalinter
install:
- make -B install-dep vendor install-gometalinter
script: make gometalinter
env:
- GOARCH=amd64
- stage: deploy
script: skip
go: "1.10"
env:
- GOARCH=amd64
deploy:
- provider: script
script: pip install awscli --user `whoami` && make -B prerelease
skip_cleanup: true
on:
all_branches: true
repo: pilosa/pilosa
tags: false
stages:
- metalinter
- test
- deploy
matrix:
# Excluding or allowing failures on non-primary matrix configurations due to long running times.
fast_finish: true
allow_failures:
- go: master
notifications:
slack:
secure: "SceWannxoGzeSu9PlEhl6icQFGuTmwax870k20nB2ZGYLjo77UEcwYoFwWvFsdYPa/HCo3JorMTYvMJ15VDJcnKEfzDr+kyXbHWBzUumclIOU/Im3ArEN6waQgyGbbWUQhvJjy4ATaxiOlmCyDV+KhKC9P3+WB33/OQtM3ngjAdTXYHAkfEcpeoOP75um+KsQgbi+hlnqfZdgDa6yIkFjaS3KZEJW1vmcOYYzNsXOA1Ip8j1NY6AjjWZlQorZJ/SYFqdhIv8ST3+a6cQk12u3t6TwZdcr3wmm1qmiW/SaK7UesWlT/YfElIuK8BBq9w1oZHxNKoAmLWTOe7MMisdItmtwgA14eMGl1rvNFlVf9sjsxs4AAzFvSZBZdDfx9XeLCBU5I2WUc/PKUgNQBPMVChxA7gEhtZLndsDdye7LsZASD2yYqjlVlgoZpzRexee/cJgCqUcNKDBHF39ZJYxV4KtZ0prjcSnVmLvuapplzTV4LZ+LyFapCyhiuM/oMJvxgmd7jTtFb5e5EkaHBPN1XwQWZw87yCjKsunTlTe1f1a5qoH/xvJHNpqE/jxOHU3DTLDgTxhb+FwC1Qj9a8bp+UYLw5F4P46ZnHlBGc2O74klv17EqvUMn3JhzASUtyxLGOgJulJ+o83rxJvhSiWt3GQIfkExVPzmz11641ElJI="

View file

@ -22,7 +22,7 @@ If you want to help but you aren't sure where to start, check out our [github la
### Development Environment
- Ensure you have a recent version of [Go](https://golang.org/doc/install) installed. Pilosa generally supports the current and previous minor versions; check our [travis file](../master/.travis.yml) for the most up-to-date information.
- Ensure you have a recent version of [Go](https://golang.org/doc/install) installed. Pilosa generally supports the current and previous minor versions; check our [CircleCI config file](../master/.circleci/config.yml) for the most up-to-date information.
- Make sure `$GOPATH` environment variable points to your Go working directory and `$PATH` incudes `$GOPATH/bin`, as described [here](https://golang.org/doc/code.html#GOPATH).

26
Dockerfile-clustertests Normal file
View file

@ -0,0 +1,26 @@
# This Dockerfile is used for cluster testing - it produces a much larger image
# and includes all of Go as well as some utilities.
FROM golang:1.11
LABEL maintainer "dev@pilosa.com"
COPY . /go/src/github.com/pilosa/pilosa/
RUN cd /go/src/github.com/pilosa/pilosa \
&& CGO_ENABLED=0 make install-dep install FLAGS="-a"
# download pumba for fault injection
ADD https://github.com/alexei-led/pumba/releases/download/0.6.0/pumba_linux_amd64 /pumba
RUN chmod +x /pumba
RUN cp /go/bin/pilosa /pilosa
COPY LICENSE /LICENSE
COPY NOTICE /NOTICE
EXPOSE 10101
VOLUME /data
ENTRYPOINT ["bash", "-c"]
CMD ["/pilosa", "server", "--data-dir", "/data", "--bind", "http://0.0.0.0:10101"]

11
Gopkg.lock generated
View file

@ -189,6 +189,15 @@
revision = "c01d1270ff3e442a8a57cddc1c92dc1138598194"
version = "v1.2.0"
[[projects]]
name = "github.com/pilosa/go-pilosa"
packages = [
".",
"gopilosa_pbuf"
]
revision = "4e7807f5ad779407936744057cd17332046b6c3c"
version = "v1.1.0"
[[projects]]
name = "github.com/pkg/errors"
packages = ["."]
@ -324,6 +333,6 @@
[solve-meta]
analyzer-name = "dep"
analyzer-version = 1
inputs-digest = "8290156ce8b4066c46ab83d743f4c81df0a17e148415bb1ee8409a51ac4c3ba4"
inputs-digest = "be318fa4f2a72e7e849b2faff1ce9300deaaf2d24cf76100f4b538abed86295b"
solver-name = "gps-cdcl"
solver-version = 1

View file

@ -68,6 +68,21 @@ release: check-clean
$(MAKE) release-build GOOS=linux GOARCH=386
$(MAKE) release-build GOOS=linux GOARCH=386 ENTERPRISE=1
# Run cluster integration tests using docker. Requires docker daemon to be
# running. This will catch changes to internal/clustertests/*.go, but if you
# make changes to Pilosa, you'll want to run clustertests-build to rebuild the
# pilosa image.
clustertests:
docker-compose -f internal/clustertests/docker-compose.yml down
docker-compose -f internal/clustertests/docker-compose.yml build client1
docker-compose -f internal/clustertests/docker-compose.yml up --exit-code-from=client1
# Like clustertests, but rebuilds all images.
clustertests-build:
docker-compose -f internal/clustertests/docker-compose.yml down
docker-compose -f internal/clustertests/docker-compose.yml up --exit-code-from=client1 --build
# Create prerelease builds
prerelease: vendor
$(MAKE) release-build GOOS=linux GOARCH=amd64 VERSION_ID=$$\(BRANCH_ID\)

View file

@ -4,7 +4,7 @@
</a>
</p>
[![Build Status](https://travis-ci.org/pilosa/pilosa.svg?branch=master)](https://travis-ci.org/pilosa/pilosa)
[![CircleCI](https://circleci.com/gh/pilosa/pilosa/tree/master.svg?style=shield)](https://circleci.com/gh/pilosa/pilosa/tree/master)
[![GoDoc](https://godoc.org/github.com/pilosa/pilosa?status.svg)](https://godoc.org/github.com/pilosa/pilosa)
[![Go Report Card](https://goreportcard.com/badge/github.com/pilosa/pilosa)](https://goreportcard.com/report/github.com/pilosa/pilosa)
[![license](https://img.shields.io/github/license/pilosa/pilosa.svg)](https://github.com/pilosa/pilosa/blob/master/LICENSE)

9
api.go
View file

@ -28,6 +28,7 @@ import (
"github.com/pilosa/pilosa/pql"
"github.com/pilosa/pilosa/roaring"
"github.com/pilosa/pilosa/stats"
"github.com/pkg/errors"
"golang.org/x/sync/errgroup"
)
@ -268,6 +269,10 @@ func setUpImportOptions(opts ...ImportOption) (*ImportOptions, error) {
// of the rows in this shard of this field concatenated together in one long
// bitmap.
func (api *API) ImportRoaring(ctx context.Context, indexName, fieldName string, shard uint64, remote bool, data []byte, opts ...ImportOption) (err error) {
if len(data) == 0 {
return errors.New("no data to import")
}
if err = api.validate(apiField); err != nil {
return errors.Wrap(err, "validating api method")
}
@ -926,7 +931,7 @@ func (api *API) AvailableShardsByIndex(_ context.Context) map[string]*roaring.Bi
// StatsWithTags returns an instance of whatever implementation of StatsClient
// pilosa is using with the given tags.
func (api *API) StatsWithTags(tags []string) StatsClient {
func (api *API) StatsWithTags(tags []string) stats.StatsClient {
if api.holder == nil || api.cluster == nil {
return nil
}
@ -952,7 +957,7 @@ func (api *API) validateShardOwnership(indexName string, shard uint64) error {
}
func (api *API) indexField(indexName string, fieldName string, shard uint64) (*Index, *Field, error) {
api.server.logger.Printf("importing: %v %v %v", indexName, fieldName, shard)
api.server.logger.Debugf("importing: %v %v %v", indexName, fieldName, shard)
// Find the Index.
index := api.holder.Index(indexName)

View file

@ -4,9 +4,9 @@ package pilosa
import "strconv"
const _apiMethod_name = "apiClusterMessageapiCreateFieldapiCreateIndexapiDeleteFieldapiDeleteIndexapiDeleteViewapiExportCSVapiFragmentBlockDataapiFragmentBlocksapiFieldapiFieldAttrDiffapiImportapiImportValueapiIndexapiIndexAttrDiffapiQueryapiRecalculateCachesapiRemoveNodeapiResizeAbortapiSetCoordinatorapiShardNodesapiViews"
const _apiMethod_name = "apiClusterMessageapiCreateFieldapiCreateIndexapiDeleteFieldapiDeleteAvailableShardapiDeleteIndexapiDeleteViewapiExportCSVapiFragmentBlockDataapiFragmentBlocksapiFieldapiFieldAttrDiffapiImportapiImportValueapiIndexapiIndexAttrDiffapiQueryapiRecalculateCachesapiRemoveNodeapiResizeAbortapiSetCoordinatorapiShardNodesapiViews"
var _apiMethod_index = [...]uint16{0, 17, 31, 45, 59, 73, 86, 98, 118, 135, 143, 159, 168, 182, 190, 206, 214, 234, 247, 261, 278, 291, 299}
var _apiMethod_index = [...]uint16{0, 17, 31, 45, 59, 82, 96, 109, 121, 141, 158, 166, 182, 191, 205, 213, 229, 237, 257, 270, 284, 301, 314, 322}
func (i apiMethod) String() string {
if i < 0 || i >= apiMethod(len(_apiMethod_index)-1) {

View file

@ -23,6 +23,7 @@ import (
"time"
"github.com/pilosa/pilosa/lru"
"github.com/pilosa/pilosa/stats"
)
const (
@ -50,14 +51,14 @@ type cache interface {
Top() []bitmapPair
// SetStats defines the stats client used in the cache.
SetStats(s StatsClient)
SetStats(s stats.StatsClient)
}
// lruCache represents a least recently used Cache implementation.
type lruCache struct {
cache *lru.Cache
counts map[uint64]uint64
stats StatsClient
stats stats.StatsClient
}
// newLRUCache returns a new instance of LRUCache.
@ -65,7 +66,7 @@ func newLRUCache(maxEntries uint32) *lruCache {
c := &lruCache{
cache: lru.New(int(maxEntries)),
counts: make(map[uint64]uint64),
stats: NopStatsClient,
stats: stats.NopStatsClient,
}
c.cache.OnEvicted = c.onEvicted
return c
@ -122,7 +123,7 @@ func (c *lruCache) Top() []bitmapPair {
}
// SetStats defines the stats client used in the cache.
func (c *lruCache) SetStats(s StatsClient) {
func (c *lruCache) SetStats(s stats.StatsClient) {
c.stats = s
}
@ -150,7 +151,7 @@ type rankCache struct {
// thresholdValue is the value of the last item in the cache
thresholdValue uint64
stats StatsClient
stats stats.StatsClient
}
// NewRankCache returns a new instance of RankCache.
@ -159,7 +160,7 @@ func NewRankCache(maxEntries uint32) *rankCache {
maxEntries: maxEntries,
thresholdBuffer: int(thresholdFactor * float64(maxEntries)),
entries: make(map[uint64]uint64),
stats: NopStatsClient,
stats: stats.NopStatsClient,
}
}
@ -279,7 +280,7 @@ func (c *rankCache) recalculate() {
}
// SetStats defines the stats client used in the cache.
func (c *rankCache) SetStats(s StatsClient) {
func (c *rankCache) SetStats(s stats.StatsClient) {
c.stats = s
}
@ -458,12 +459,12 @@ func (s *simpleCache) Add(id uint64, b *Row) {
// nopCache represents a no-op Cache implementation.
type nopCache struct {
stats StatsClient
stats stats.StatsClient
}
// Ensure NopCache implements Cache.
var globalNopCache cache = nopCache{
stats: NopStatsClient,
stats: stats.NopStatsClient,
}
func (c nopCache) Add(uint64, uint64) {}
@ -471,10 +472,10 @@ func (c nopCache) BulkAdd(uint64, uint64) {}
func (c nopCache) Get(uint64) uint64 { return 0 }
func (c nopCache) IDs() []uint64 { return []uint64{} }
func (c nopCache) Invalidate() {}
func (c nopCache) Len() int { return 0 }
func (c nopCache) Recalculate() {}
func (c nopCache) SetStats(StatsClient) {}
func (c nopCache) Invalidate() {}
func (c nopCache) Len() int { return 0 }
func (c nopCache) Recalculate() {}
func (c nopCache) SetStats(stats.StatsClient) {}
func (c nopCache) Top() []bitmapPair {
return []bitmapPair{}

View file

@ -31,6 +31,7 @@ import (
"github.com/gogo/protobuf/proto"
"github.com/pilosa/pilosa/internal"
"github.com/pilosa/pilosa/logger"
"github.com/pilosa/pilosa/roaring"
"github.com/pkg/errors"
uuid "github.com/satori/go.uuid"
@ -65,10 +66,11 @@ type Node struct {
ID string `json:"id"`
URI URI `json:"uri"`
IsCoordinator bool `json:"isCoordinator"`
State string `json:"state"`
}
func (n Node) String() string {
return fmt.Sprintf("Node: %s", n.ID)
return fmt.Sprintf("Node:%s:%s:%s", n.URI, n.State, n.ID[:6])
}
// Nodes represents a list of nodes.
@ -215,7 +217,7 @@ type cluster struct { // nolint: maligned
wg sync.WaitGroup
closing chan struct{}
logger Logger
logger logger.Logger
InternalClient InternalClient
}
@ -234,7 +236,7 @@ func newCluster() *cluster {
InternalClient: newNopInternalClient(),
logger: NopLogger,
logger: logger.NopLogger,
}
}
@ -345,8 +347,6 @@ func (c *cluster) unprotectedUpdateCoordinator(n *Node) bool {
// addNode adds a node to the Cluster and updates and saves the
// new topology. unprotected.
func (c *cluster) addNode(node *Node) error {
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 {
c.Coordinator = node.ID
@ -456,7 +456,19 @@ func (c *cluster) unprotectedSetState(state string) {
}
}
func (c *cluster) setMyNodeState(state string) {
c.mu.Lock()
defer c.mu.Unlock()
c.Node.State = state
for i, n := range c.nodes {
if n.ID == c.Node.ID {
c.nodes[i].State = state
}
}
}
func (c *cluster) setNodeState(state string) error { // nolint: unparam
c.setMyNodeState(state)
if c.isCoordinator() {
return c.receiveNodeState(c.Node.ID, state)
}
@ -467,7 +479,7 @@ func (c *cluster) setNodeState(state string) error { // nolint: unparam
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)
}
@ -486,11 +498,23 @@ func (c *cluster) receiveNodeState(nodeID string, state string) error {
}
c.Topology.mu.Lock()
c.Topology.nodeStates[nodeID] = state
changed := false
if c.Topology.nodeStates[nodeID] != state {
changed = true
c.Topology.nodeStates[nodeID] = state
for i, n := range c.nodes {
if n.ID == nodeID {
c.nodes[i].State = state
}
}
}
c.Topology.mu.Unlock()
c.logger.Printf("received state %s (%s)", state, nodeID)
return c.unprotectedSetStateAndBroadcast(c.determineClusterState())
if changed {
return c.unprotectedSetStateAndBroadcast(c.determineClusterState())
}
return nil
}
// determineClusterState is unprotected.
@ -932,7 +956,6 @@ func (c *cluster) waitForStarted() error {
<-c.joining
c.logger.Printf("joining has completed")
}
return nil
}
@ -945,7 +968,6 @@ func (c *cluster) close() error {
}
func (c *cluster) markAsJoined() {
c.logger.Printf("mark node as joined (received coordinator update)")
if !c.joined {
c.joined = true
close(c.joining)
@ -1043,8 +1065,8 @@ func (c *cluster) unprotectedSetStateAndBroadcast(state string) error {
return nil
}
// Broadcast cluster status changes to the cluster.
c.logger.Printf("broadcasting ClusterStatus: %s", state)
return c.broadcaster.SendSync(c.unprotectedStatus()) // TODO fix c.Status
status := c.unprotectedStatus()
return c.broadcaster.SendSync(status) // TODO fix c.Status
}
@ -1220,14 +1242,14 @@ func (c *cluster) followResizeInstruction(instr *ResizeInstruction) error {
return errors.Wrap(err, "merging cluster status")
}
c.logger.Printf("MergeClusterStatus done, start goroutine")
c.logger.Printf("done MergeClusterStatus, 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.Recv()
// Prepare the return message.
complete := &ResizeInstructionComplete{
@ -1240,7 +1262,7 @@ func (c *cluster) followResizeInstruction(instr *ResizeInstruction) error {
if err := func() error {
// Sync the schema received in the resize instruction.
c.logger.Printf("Holder ApplySchema")
c.logger.Debugf("holder applySchema")
if err := c.holder.applySchema(instr.Schema); err != nil {
return errors.Wrap(err, "applying schema")
}
@ -1354,7 +1376,7 @@ type resizeJob struct {
mu sync.RWMutex
state string
Logger Logger
Logger logger.Logger
}
// newResizeJob returns a new instance of resizeJob.
@ -1386,7 +1408,7 @@ func newResizeJob(existingNodes []*Node, node *Node, action string) *resizeJob {
IDs: ids,
action: action,
result: make(chan string),
Logger: NopLogger,
Logger: logger.NopLogger,
}
}
@ -1625,17 +1647,17 @@ func (c *cluster) ReceiveEvent(e *NodeEvent) (err error) {
switch e.Event {
case NodeJoin:
c.logger.Printf("received NodeJoin event: %v", e)
c.logger.Debugf("nodeJoin of %s on %s", e.Node.URI, c.Node.URI)
// Ignore the event if this is not the coordinator.
if !c.isCoordinator() {
return nil
}
return c.nodeJoin(e.Node)
case NodeLeave:
c.logger.Printf("received node leave on %s: %s, uri: %v", c.Node, e.Node, e.Node.URI)
c.mu.Lock()
defer c.mu.Unlock()
if c.unprotectedIsCoordinator() {
c.logger.Printf("received node leave: %v", e.Node)
// if removeNodeBasicSorted succeeds, that means that the node was
// not already removed by a removeNode request. We treat this as the
// host being temporarily unavailable, and expect it to come back
@ -1647,7 +1669,6 @@ func (c *cluster) ReceiveEvent(e *NodeEvent) (err error) {
err = c.unprotectedSetStateAndBroadcast(c.determineClusterState())
}
}
c.logger.Printf("finished node leave on %s: %s, uri: %v", c.Node, e.Node, e.Node.URI)
case NodeUpdate:
c.logger.Printf("received node update event: id: %v, string: %v, uri: %v", e.Node.ID, e.Node.String(), e.Node.URI)
// NodeUpdate is intentionally not implemented.
@ -1660,6 +1681,7 @@ func (c *cluster) ReceiveEvent(e *NodeEvent) (err error) {
func (c *cluster) nodeJoin(node *Node) error {
c.mu.Lock()
defer c.mu.Unlock()
c.logger.Printf("node join event on coordinator, node: %s, id: %s", node.URI, node.ID)
if c.needTopologyAgreement() {
// A host that is not part of the topology can't be added to the STARTING cluster.
if !c.Topology.ContainsID(node.ID) {
@ -1688,11 +1710,10 @@ func (c *cluster) nodeJoin(node *Node) error {
if c.haveTopologyAgreement() && c.allNodesReady() {
return c.unprotectedSetStateAndBroadcast(ClusterStateNormal)
} else {
// Send the status to the remote node. This lets the remote node
// know that it can proceed with opening its Holder.
return c.sendTo(node, c.unprotectedStatus())
}
// Send the status to the remote node. This lets the remote node
// know that it can proceed with opening its Holder.
return c.sendTo(node, c.unprotectedStatus())
}
// If the cluster already contains the node, just send it the cluster status.
@ -1700,7 +1721,7 @@ func (c *cluster) nodeJoin(node *Node) error {
// the cluster.
if cnode := c.unprotectedNodeByID(node.ID); cnode != nil {
if cnode.URI != node.URI {
c.logger.Printf("Node: %v changed URI from %s to %s", cnode.ID, cnode.URI, node.URI)
c.logger.Printf("node: %v changed URI from %s to %s", cnode.ID, cnode.URI, node.URI)
cnode.URI = node.URI
}
return c.unprotectedSetStateAndBroadcast(c.determineClusterState())
@ -1796,6 +1817,10 @@ func (c *cluster) mergeClusterStatus(cs *ClusterStatus) error {
// Add all nodes from the coordinator.
for _, node := range officialNodes {
if node.ID == c.Node.ID && node.State != c.Node.State {
c.logger.Printf("mismatched state in mergeClusterStatus got %v have %v", node.State, c.Node.State)
go c.setNodeState(c.Node.State)
}
if err := c.addNode(node); err != nil {
return errors.Wrap(err, "adding node")
}

View file

@ -185,8 +185,8 @@ func TestFragSources(t *testing.T) {
"node0": {},
"node1": {},
"node2": {
{&Node{"node0", URI{"http", "host0", 10101}, false}, "i", "f", "standard", uint64(0)},
{&Node{"node1", URI{"http", "host1", 10101}, false}, "i", "f", "standard", uint64(2)},
{&Node{ID: "node0", URI: URI{"http", "host0", 10101}, IsCoordinator: false}, "i", "f", "standard", uint64(0)},
{&Node{ID: "node1", URI: URI{"http", "host1", 10101}, IsCoordinator: false}, "i", "f", "standard", uint64(2)},
},
},
err: "",
@ -197,11 +197,11 @@ func TestFragSources(t *testing.T) {
idx: idx,
expected: map[string][]*ResizeSource{
"node0": {
{&Node{"node1", URI{"http", "host1", 10101}, false}, "i", "f", "standard", uint64(1)},
{&Node{ID: "node1", URI: URI{"http", "host1", 10101}, IsCoordinator: false}, "i", "f", "standard", uint64(1)},
},
"node1": {
{&Node{"node0", URI{"http", "host0", 10101}, false}, "i", "f", "standard", uint64(0)},
{&Node{"node0", URI{"http", "host0", 10101}, false}, "i", "f", "standard", uint64(2)},
{&Node{ID: "node0", URI: URI{"http", "host0", 10101}, IsCoordinator: false}, "i", "f", "standard", uint64(0)},
{&Node{ID: "node0", URI: URI{"http", "host0", 10101}, IsCoordinator: false}, "i", "f", "standard", uint64(2)},
},
},
err: "",
@ -212,11 +212,11 @@ func TestFragSources(t *testing.T) {
idx: idx,
expected: map[string][]*ResizeSource{
"node0": {
{&Node{"node2", URI{"http", "host2", 10101}, false}, "i", "f", "standard", uint64(0)},
{&Node{"node2", URI{"http", "host2", 10101}, false}, "i", "f", "standard", uint64(2)},
{&Node{ID: "node2", URI: URI{"http", "host2", 10101}, IsCoordinator: false}, "i", "f", "standard", uint64(0)},
{&Node{ID: "node2", URI: URI{"http", "host2", 10101}, IsCoordinator: false}, "i", "f", "standard", uint64(2)},
},
"node1": {
{&Node{"node0", URI{"http", "host0", 10101}, false}, "i", "f", "standard", uint64(3)},
{&Node{ID: "node0", URI: URI{"http", "host0", 10101}, IsCoordinator: false}, "i", "f", "standard", uint64(3)},
},
"node2": {},
},

View file

@ -19,10 +19,8 @@ import (
"io"
"github.com/pilosa/pilosa"
"github.com/spf13/cobra"
"github.com/pilosa/pilosa/ctl"
"github.com/spf13/cobra"
)
var Importer *ctl.ImportCommand
@ -55,6 +53,7 @@ omitted. If it is present then its format should be YYYY-MM-DDTHH:MM.
flags.StringVarP(&Importer.Field, "field", "f", "", "Field to import into.")
flags.BoolVar(&Importer.IndexOptions.Keys, "index-keys", false, "Specify keys=true when creating an index")
flags.BoolVar(&Importer.FieldOptions.Keys, "field-keys", false, "Specify keys=true when creating a field")
flags.StringVar(&Importer.FieldOptions.Type, "field-type", "", "Specify the field type when creating a field. One of: set, int, time, bool, mutex")
flags.Int64Var(&Importer.FieldOptions.Min, "field-min", 0, "Specify the minimum for an int field on creation")
flags.Int64Var(&Importer.FieldOptions.Max, "field-max", 0, "Specify the maximum for an int field on creation")
flags.StringVar(&Importer.FieldOptions.CacheType, "field-cache-type", pilosa.CacheTypeRanked, "Specify the cache type for a set field on creation. One of: none, lru, ranked")

View file

@ -99,13 +99,15 @@ func (cmd *ImportCommand) Run(ctx context.Context) error {
cmd.client = client
if cmd.CreateSchema {
// set the correct type for the field
if cmd.FieldOptions.TimeQuantum != "" {
cmd.FieldOptions.Type = "time"
} else if cmd.FieldOptions.Min != 0 || cmd.FieldOptions.Max != 0 {
cmd.FieldOptions.Type = "int"
} else {
cmd.FieldOptions.Type = "set"
if cmd.FieldOptions.Type == "" {
// set the correct type for the field
if cmd.FieldOptions.TimeQuantum != "" {
cmd.FieldOptions.Type = "time"
} else if cmd.FieldOptions.Min != 0 || cmd.FieldOptions.Max != 0 {
cmd.FieldOptions.Type = "int"
} else {
cmd.FieldOptions.Type = "set"
}
}
err := cmd.ensureSchema(ctx)
if err != nil {

View file

@ -24,6 +24,7 @@ import (
"sync"
"time"
"github.com/pilosa/pilosa/logger"
"github.com/pkg/errors"
)
@ -51,7 +52,7 @@ type diagnosticsCollector struct {
client *http.Client
Logger Logger
Logger logger.Logger
server *Server
}
@ -65,7 +66,7 @@ func newDiagnosticsCollector(host string) *diagnosticsCollector { // nolint: unp
start: time.Now(),
client: &http.Client{Timeout: 10 * time.Second},
metrics: make(map[string]interface{}),
Logger: NopLogger,
Logger: logger.NopLogger,
}
}

View file

@ -506,6 +506,7 @@ func encodeNode(n *pilosa.Node) *internal.Node {
ID: n.ID,
URI: encodeURI(n.URI),
IsCoordinator: n.IsCoordinator,
State: n.State,
}
}
@ -761,6 +762,7 @@ func decodeNode(node *internal.Node, m *pilosa.Node) {
m.ID = node.ID
decodeURI(node.URI, &m.URI)
m.IsCoordinator = node.IsCoordinator
m.State = node.State
}
func decodeURI(i *internal.URI, m *pilosa.URI) {

View file

@ -1664,7 +1664,7 @@ func (e *executor) executeSetRowShard(ctx context.Context, index string, c *pql.
if err != nil {
return false, errors.Wrap(err, "creating view")
}
fragment, err = view.createFragmentIfNotExists(shard)
fragment, err = view.CreateFragmentIfNotExists(shard)
if err != nil {
return false, errors.Wrapf(err, "creating fragment: %d", shard)
}
@ -2055,7 +2055,7 @@ func (e *executor) mapReduce(ctx context.Context, index string, shards []uint64,
if !opt.Remote {
nodes = Nodes(e.Cluster.nodes).Clone()
} else {
nodes = []*Node{e.Cluster.unprotectedNodeByID(e.Node.ID)}
nodes = []*Node{e.Cluster.nodeByID(e.Node.ID)}
}
// Start mapping across all primary owners.

View file

@ -28,8 +28,10 @@ import (
"github.com/gogo/protobuf/proto"
"github.com/pilosa/pilosa/internal"
"github.com/pilosa/pilosa/logger"
"github.com/pilosa/pilosa/pql"
"github.com/pilosa/pilosa/roaring"
"github.com/pilosa/pilosa/stats"
"github.com/pkg/errors"
)
@ -69,7 +71,7 @@ type Field struct {
rowAttrStore AttrStore
broadcaster broadcaster
Stats StatsClient
Stats stats.StatsClient
// Field options.
options FieldOptions
@ -79,7 +81,7 @@ type Field struct {
// Shards with data on any node in the cluster, according to this node.
remoteAvailableShards *roaring.Bitmap
logger Logger
logger logger.Logger
}
// FieldOption is a functional option type for pilosa.fieldOptions.
@ -199,13 +201,13 @@ func newField(path, index, name string, opts FieldOption) (*Field, error) {
rowAttrStore: nopStore,
broadcaster: NopBroadcaster,
Stats: NopStatsClient,
Stats: stats.NopStatsClient,
options: applyDefaultOptions(fo),
remoteAvailableShards: roaring.NewBitmap(),
logger: NopLogger,
logger: logger.NopLogger,
}
return f, nil
}

View file

@ -25,6 +25,7 @@ import (
"hash"
"io"
"io/ioutil"
"math"
"os"
"sort"
"sync"
@ -33,13 +34,12 @@ import (
"unsafe"
"github.com/cespare/xxhash"
"math"
"github.com/gogo/protobuf/proto"
"github.com/pilosa/pilosa/internal"
"github.com/pilosa/pilosa/logger"
"github.com/pilosa/pilosa/pql"
"github.com/pilosa/pilosa/roaring"
"github.com/pilosa/pilosa/stats"
"github.com/pkg/errors"
)
@ -118,7 +118,7 @@ type fragment struct {
MaxOpN int
// Logger used for out-of-band log entries.
Logger Logger
Logger logger.Logger
// Row attribute storage.
// This is set by the parent field unless overridden for testing.
@ -128,7 +128,7 @@ type fragment struct {
// existing value (to clear) prior to setting a new value.
mutexVector vector
stats StatsClient
stats stats.StatsClient
}
// newFragment returns a new instance of Fragment.
@ -142,10 +142,10 @@ func newFragment(path, index, field, view string, shard uint64) *fragment {
CacheType: DefaultCacheType,
CacheSize: DefaultCacheSize,
Logger: NopLogger,
Logger: logger.NopLogger,
MaxOpN: defaultFragmentMaxOpN,
stats: NopStatsClient,
stats: stats.NopStatsClient,
}
}
@ -1492,9 +1492,6 @@ func (f *fragment) bulkImportStandard(rowIDs, columnIDs []uint64, options *Impor
lastRowID = rowID
rowSet[rowID] = struct{}{}
}
// Invalidate block checksum.
delete(f.checksums, int(rowID/HashBlockSize))
}
f.mu.Lock()
@ -1518,6 +1515,9 @@ func (f *fragment) bulkImportStandard(rowIDs, columnIDs []uint64, options *Impor
// Update cache counts for all affected rows.
for rowID := range rowSet {
// Invalidate block checksum.
delete(f.checksums, int(rowID/HashBlockSize))
n := results.CountRange(rowID*ShardWidth, (rowID+1)*ShardWidth)
f.cache.BulkAdd(rowID, n)
}
@ -1720,7 +1720,7 @@ func (f *fragment) Snapshot() error {
defer f.mu.Unlock()
return f.snapshot()
}
func track(start time.Time, message string, stats StatsClient, logger Logger) {
func track(start time.Time, message string, stats stats.StatsClient, logger logger.Logger) {
elapsed := time.Since(start)
logger.Printf("%s took %s", message, elapsed)
stats.Histogram("snapshot", elapsed.Seconds(), 1.0)
@ -1734,7 +1734,6 @@ func (f *fragment) snapshot() error {
// f.mu must be locked when calling it.
func unprotectedWriteToFragment(f *fragment, bm *roaring.Bitmap) error { // nolint: interfacer
f.Logger.Printf("fragment: snapshotting %s/%s/%s/%d", f.index, f.field, f.view, f.shard)
completeMessage := fmt.Sprintf("fragment: snapshot complete %s/%s/%s/%d", f.index, f.field, f.view, f.shard)
start := time.Now()
defer track(start, completeMessage, f.stats, f.Logger)

View file

@ -25,6 +25,8 @@ import (
"testing"
"testing/quick"
"golang.org/x/sync/errgroup"
"github.com/davecgh/go-spew/spew"
"github.com/pilosa/pilosa/pql"
"github.com/pilosa/pilosa/roaring"
@ -1399,6 +1401,21 @@ func TestFragment_ImportSet(t *testing.T) {
}
}
func TestFragment_ConcurrentImport(t *testing.T) {
t.Run("bulkImportStandard", func(t *testing.T) {
f := mustOpenFragment("i", "f", viewStandard, 0, "")
defer f.Close()
eg := errgroup.Group{}
eg.Go(func() error { return f.bulkImportStandard([]uint64{1, 2}, []uint64{1, 2}, &ImportOptions{}) })
eg.Go(func() error { return f.bulkImportStandard([]uint64{3, 4}, []uint64{3, 4}, &ImportOptions{}) })
err := eg.Wait()
if err != nil {
t.Fatalf("importing data to fragment: %v", err)
}
})
}
// Ensure a fragment can import mutually exclusive values.
func TestFragment_ImportMutex(t *testing.T) {
tests := []struct {

View file

@ -18,6 +18,7 @@ import (
"bytes"
"context"
"fmt"
"io"
"io/ioutil"
"log"
"net"
@ -28,6 +29,7 @@ import (
"github.com/hashicorp/memberlist"
"github.com/pilosa/pilosa"
"github.com/pilosa/pilosa/logger"
"github.com/pilosa/pilosa/roaring"
"github.com/pilosa/pilosa/toml"
"github.com/pkg/errors"
@ -46,9 +48,10 @@ type memberSet struct {
papi *pilosa.API
config *config
Logger pilosa.Logger
Logger logger.Logger
logger *log.Logger
logOutput io.Writer
transport *Transport
eventReceiver *eventReceiver
@ -156,12 +159,19 @@ func WithLogger(logger *log.Logger) memberSetOption {
}
}
func WithLogOutput(o io.Writer) memberSetOption {
return func(g *memberSet) error {
g.logOutput = o
return nil
}
}
// 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 := &memberSet{
papi: api,
Logger: pilosa.NopLogger,
Logger: logger.NopLogger,
}
// options
@ -220,7 +230,11 @@ func NewMemberSet(cfg Config, api *pilosa.API, options ...memberSetOption) (*mem
conf.Delegate = g
conf.SecretKey = gossipKey
conf.Events = ger
conf.Logger = g.logger
if g.logOutput != nil {
conf.LogOutput = g.logOutput
} else {
conf.Logger = g.logger
}
g.config = &config{
memberlistConfig: conf,

View file

@ -27,7 +27,9 @@ import (
"syscall"
"time"
"github.com/pilosa/pilosa/logger"
"github.com/pilosa/pilosa/roaring"
"github.com/pilosa/pilosa/stats"
"github.com/pkg/errors"
uuid "github.com/satori/go.uuid"
)
@ -55,7 +57,7 @@ type Holder struct {
NewPrimaryTranslateStore func(interface{}) TranslateStore
// opened channel is closed once Open() completes.
opened chan struct{}
opened lockedChan
broadcaster broadcaster
@ -66,7 +68,7 @@ type Holder struct {
closing chan struct{}
// Stats
Stats StatsClient
Stats stats.StatsClient
// Data directory path.
Path string
@ -74,7 +76,33 @@ type Holder struct {
// The interval at which the cached row ids are persisted to disk.
cacheFlushInterval time.Duration
Logger Logger
Logger logger.Logger
}
// lockedChan looks a little ridiculous admittedly, but exists for good reason.
// The channel within is used (for example) to signal to other goroutines when
// the Holder has finished opening (via closing the channel). However, it is
// possible for the holder to be closed and then reopened, but a channel which
// is closed cannot be re-opened. We must create a new channel - this creates a
// data race with any goroutine which might be accessing the channel. To ensure
// that there is no data race on the value of the channel itself, we wrap any
// operation on it with an RWMutex so that we can guarantee that nothing is
// trying to listen on it when it gets swapped.
type lockedChan struct {
ch chan struct{}
mu sync.RWMutex
}
func (lc *lockedChan) Close() {
lc.mu.RLock()
close(lc.ch)
lc.mu.RUnlock()
}
func (lc *lockedChan) Recv() {
lc.mu.RLock()
<-lc.ch
lc.mu.RUnlock()
}
// NewHolder returns a new instance of Holder.
@ -83,19 +111,19 @@ func NewHolder() *Holder {
indexes: make(map[string]*Index),
closing: make(chan struct{}),
opened: make(chan struct{}),
opened: lockedChan{ch: make(chan struct{})},
translateFile: NewTranslateFile(),
NewPrimaryTranslateStore: newNopTranslateStore,
broadcaster: NopBroadcaster,
Stats: NopStatsClient,
Stats: stats.NopStatsClient,
NewAttrStore: newNopAttrStore,
cacheFlushInterval: defaultCacheFlushInterval,
Logger: NopLogger,
Logger: logger.NopLogger,
}
}
@ -157,7 +185,7 @@ func (h *Holder) Open() error {
h.Stats.Open()
close(h.opened)
h.opened.Close()
return nil
}
@ -182,7 +210,9 @@ func (h *Holder) Close() error {
}
// Reset opened in case Holder needs to be reopened.
h.opened = make(chan struct{})
h.opened.mu.Lock()
h.opened.ch = make(chan struct{})
h.opened.mu.Unlock()
return nil
}
@ -472,7 +502,7 @@ func (h *Holder) flushCaches() {
}
if err := fragment.FlushCache(); err != nil {
h.Logger.Printf("error flushing cache: err=%s, path=%s", err, fragment.cachePath())
h.Logger.Printf("ERROR flushing cache: err=%s, path=%s", err, fragment.cachePath())
}
}
}
@ -533,7 +563,7 @@ func (h *Holder) setFileLimit() {
h.Logger.Printf("ERROR checking open file limit: %s", err)
} else {
if oldLimit.Cur < fileLimit {
h.Logger.Printf("WARNING: Tried to set open file limit to %d, but it is %d. You may consider running \"sudo ulimit -n %d\" before starting Pilosa to avoid \"too many open files\" error. See https://www.pilosa.com/docs/administration/#open-file-limits for more information.", fileLimit, oldLimit.Cur, fileLimit)
h.Logger.Printf("WARNING: Tried to set open file limit to %d, but it is %d. You may consider running \"sudo ulimit -n %d\" before starting Pilosa to avoid \"too many open files\" error. See https://www.pilosa.com/docs/latest/administration/#open-file-limits for more information.", fileLimit, oldLimit.Cur, fileLimit)
}
}
}
@ -605,7 +635,7 @@ type holderSyncer struct {
Cluster *cluster
// Stats
Stats StatsClient
Stats stats.StatsClient
// Signals that the sync should stop.
Closing <-chan struct{}

View file

@ -91,9 +91,7 @@ func (c *InternalClient) maxShardByIndex(ctx context.Context) (map[string]uint64
defer resp.Body.Close()
var rsp getShardsMaxResponse
if resp.StatusCode != http.StatusOK {
return nil, fmt.Errorf("http: status=%d", resp.StatusCode)
} else if err := json.NewDecoder(resp.Body).Decode(&rsp); err != nil {
if err := json.NewDecoder(resp.Body).Decode(&rsp); err != nil {
return nil, fmt.Errorf("json decode: %s", err)
}
@ -152,7 +150,7 @@ func (c *InternalClient) CreateIndex(ctx context.Context, index string, opt pilo
// Execute request against the host.
resp, err := c.executeRequest(req.WithContext(ctx))
if err != nil {
if resp.StatusCode == http.StatusConflict {
if resp != nil && resp.StatusCode == http.StatusConflict {
return pilosa.ErrIndexExists
}
return err
@ -258,8 +256,6 @@ func (c *InternalClient) QueryNode(ctx context.Context, uri *pilosa.URI, index s
body, err := ioutil.ReadAll(resp.Body)
if err != nil {
return nil, errors.Wrap(err, "reading")
} else if resp.StatusCode != http.StatusOK {
return nil, errors.New(string(body))
}
qresp := &pilosa.QueryResponse{}
@ -689,7 +685,7 @@ func (c *InternalClient) backupShardNode(ctx context.Context, index, field strin
// Execute request.
resp, err := c.executeRequest(req.WithContext(ctx))
if err != nil {
if resp.StatusCode == http.StatusNotFound {
if resp != nil && resp.StatusCode == http.StatusNotFound {
return nil, pilosa.ErrFragmentNotFound
}
return nil, err
@ -746,7 +742,7 @@ func (c *InternalClient) CreateFieldWithOptions(ctx context.Context, index, fiel
// Execute request against the host.
resp, err := c.executeRequest(req.WithContext(ctx))
if err != nil {
if resp.StatusCode == http.StatusConflict {
if resp != nil && resp.StatusCode == http.StatusConflict {
return pilosa.ErrFieldExists
}
return err
@ -782,7 +778,7 @@ func (c *InternalClient) FragmentBlocks(ctx context.Context, uri *pilosa.URI, in
resp, err := c.executeRequest(req.WithContext(ctx))
if err != nil {
// Return the appropriate error.
if resp.StatusCode == http.StatusNotFound {
if resp != nil && resp.StatusCode == http.StatusNotFound {
return nil, pilosa.ErrFragmentNotFound
}
return nil, err
@ -825,7 +821,7 @@ func (c *InternalClient) BlockData(ctx context.Context, uri *pilosa.URI, index,
resp, err := c.executeRequest(req.WithContext(ctx))
if err != nil {
if resp.StatusCode == http.StatusNotFound {
if resp != nil && resp.StatusCode == http.StatusNotFound {
return nil, nil, nil
}
return nil, nil, err
@ -904,7 +900,7 @@ func (c *InternalClient) RowAttrDiff(ctx context.Context, uri *pilosa.URI, index
// Execute request.
resp, err := c.executeRequest(req.WithContext(ctx))
if err != nil {
if resp.StatusCode == http.StatusNotFound {
if resp != nil && resp.StatusCode == http.StatusNotFound {
return nil, pilosa.ErrFieldNotFound
}
return nil, err

View file

@ -35,6 +35,7 @@ import (
"github.com/gorilla/handlers"
"github.com/gorilla/mux"
"github.com/pilosa/pilosa"
"github.com/pilosa/pilosa/logger"
"github.com/pkg/errors"
)
@ -42,7 +43,7 @@ import (
type Handler struct {
Handler http.Handler
logger pilosa.Logger
logger logger.Logger
// Keeps the query argument validators for each handler
validators map[string]*queryValidationSpec
@ -93,7 +94,7 @@ func OptHandlerAPI(api *pilosa.API) handlerOption {
}
}
func OptHandlerLogger(logger pilosa.Logger) handlerOption {
func OptHandlerLogger(logger logger.Logger) handlerOption {
return func(h *Handler) error {
h.logger = logger
return nil
@ -119,7 +120,7 @@ func OptHandlerCloseTimeout(d time.Duration) handlerOption {
// NewHandler returns a new instance of Handler with a default logger.
func NewHandler(opts ...handlerOption) (*Handler, error) {
handler := &Handler{
logger: pilosa.NopLogger,
logger: logger.NopLogger,
closeTimeout: time.Second * 30,
}
handler.Handler = newRouter(handler)

View file

@ -25,7 +25,9 @@ import (
"github.com/gogo/protobuf/proto"
"github.com/pilosa/pilosa/internal"
"github.com/pilosa/pilosa/logger"
"github.com/pilosa/pilosa/roaring"
"github.com/pilosa/pilosa/stats"
"github.com/pkg/errors"
)
@ -49,9 +51,9 @@ type Index struct {
columnAttrs AttrStore
broadcaster broadcaster
Stats StatsClient
Stats stats.StatsClient
logger Logger
logger logger.Logger
}
// NewIndex returns a new instance of Index.
@ -70,8 +72,8 @@ func NewIndex(path, name string) (*Index, error) {
columnAttrs: nopStore,
broadcaster: NopBroadcaster,
Stats: NopStatsClient,
logger: NopLogger,
Stats: stats.NopStatsClient,
logger: logger.NopLogger,
trackExistence: true,
}, nil
}

View file

@ -0,0 +1,3 @@
FROM ptest
COPY . /go/src/github.com/pilosa/pilosa/internal/clustertests

View file

@ -0,0 +1,81 @@
package clustertest
import (
"context"
"os"
"os/exec"
"testing"
"time"
"github.com/pilosa/pilosa"
picli "github.com/pilosa/pilosa/http"
)
func TestClusterStuff(t *testing.T) {
if os.Getenv("ENABLE_PILOSA_CLUSTER_TESTS") != "1" {
t.Skip()
}
cli, err := picli.NewInternalClient("pilosa1:10101", picli.GetHTTPClient(nil))
if err != nil {
t.Fatalf("getting client: %v", err)
}
t.Run("long pause", func(t *testing.T) {
err := cli.CreateIndex(context.Background(), "testidx", pilosa.IndexOptions{})
if err != nil {
t.Fatalf("creating index: %v", err)
}
err = cli.CreateFieldWithOptions(context.Background(), "testidx", "testf", pilosa.FieldOptions{CacheType: pilosa.CacheTypeRanked, CacheSize: 100})
if err != nil {
t.Fatalf("creating field: %v", err)
}
data := make([]pilosa.Bit, 10)
for i := 0; i < 1000; i++ {
data[i%10].RowID = 0
data[i%10].ColumnID = uint64((i/10)*pilosa.ShardWidth + i%10)
shard := uint64(i / 10)
if i%10 == 9 {
err = cli.Import(context.Background(), "testidx", "testf", shard, data)
if err != nil {
t.Fatalf("importing: %v", err)
}
}
}
r, err := cli.Query(context.Background(), "testidx", &pilosa.QueryRequest{Index: "testidx", Query: "Count(Row(testf=0))"})
if err != nil {
t.Fatalf("count querying: %v", err)
}
if r.Results[0].(uint64) != 1000 {
t.Fatalf("count after import is %d", r.Results[0].(uint64))
}
pcmd := exec.Command("/pumba", "pause", "clustertests_pilosa3_1", "--duration", "10s")
pcmd.Stdout = os.Stdout
pcmd.Stderr = os.Stderr
t.Log("pausing pilosa3 for 10s")
err = pcmd.Start()
if err != nil {
t.Fatalf("starting pumba command: %v", err)
}
err = pcmd.Wait()
if err != nil {
t.Fatalf("waiting on pumba pause cmd: %v", err)
}
// TODO change the sleep to wait for status to return to NORMAL - need support in internal client for getting status
t.Log("done with pause, waiting for stability")
time.Sleep(time.Second * 3)
t.Log("done waiting for stability")
r, err = cli.Query(context.Background(), "testidx", &pilosa.QueryRequest{Index: "testidx", Query: "Count(Row(testf=0))"})
if err != nil {
t.Fatalf("count querying: %v", err)
}
if r.Results[0].(uint64) != 1000 {
t.Fatalf("count after import is %d", r.Results[0].(uint64))
}
})
}

View file

@ -0,0 +1,55 @@
version: '2'
services:
pilosa1:
build:
context: ../..
dockerfile: Dockerfile-clustertests
image: ptest
ports:
- "33455:10101"
environment:
- PILOSA_CLUSTER_COORDINATOR=true
- PILOSA_GOSSIP_SEEDS=pilosa1:14000
networks:
- pilosanet
command:
- "/pilosa server --bind pilosa1:10101"
pilosa2:
build:
context: ../..
dockerfile: Dockerfile-clustertests
image: ptest
ports:
- "33456:10101"
environment:
- PILOSA_GOSSIP_SEEDS=pilosa1:14000
networks:
- pilosanet
command:
- "/pilosa server --bind pilosa2:10101"
pilosa3:
build:
context: ../..
dockerfile: Dockerfile-clustertests
image: ptest
ports:
- "33457:10101"
environment:
- PILOSA_GOSSIP_SEEDS=pilosa1:14000,pilosa2:14000
networks:
- pilosanet
command:
- "/pilosa server --bind pilosa3:10101"
client1:
build:
context: .
environment:
- ENABLE_PILOSA_CLUSTER_TESTS=1
networks:
- pilosanet
volumes:
- /var/run/docker.sock:/var/run/docker.sock
command:
- "go test -v -count=1 github.com/pilosa/pilosa/internal/clustertests"
networks:
pilosanet:

File diff suppressed because it is too large Load diff

View file

@ -100,6 +100,7 @@ message Node {
string ID = 1;
URI URI = 2;
bool IsCoordinator = 3;
string State = 4;
}
message NodeStateMessage {

File diff suppressed because it is too large Load diff

View file

@ -12,7 +12,7 @@
// See the License for the specific language governing permissions and
// limitations under the License.
package pilosa
package logger
import (
"io"
@ -39,7 +39,7 @@ func (n *nopLogger) Printf(format string, v ...interface{}) {}
// Debugf is a no-op implementation of the Logger Debugf method.
func (n *nopLogger) Debugf(format string, v ...interface{}) {}
// standardLogger is a basic implementation of pilosa.Logger based on log.Logger.
// standardLogger is a basic implementation of Logger based on log.Logger.
type standardLogger struct {
logger *log.Logger
}
@ -60,7 +60,7 @@ func (s *standardLogger) Logger() *log.Logger {
return s.logger
}
// verboseLogger is an implementation of pilosa.Logger which includes debug messages.
// verboseLogger is an implementation of Logger which includes debug messages.
type verboseLogger struct {
logger *log.Logger
}

View file

@ -57,7 +57,6 @@ field <- <fieldExpr / reserved> { p.addField(buffer[begin:end]) }
reserved <- ('_row' / '_col' / '_start' / '_end' / '_timestamp' / '_field')
posfield <- <fieldExpr> { p.addPosStr("_field", buffer[begin:end]) }
uint <- [1-9] [0-9]* / '0'
uintrow <- <uint>{p.addPosNum("_row", buffer[begin:end])}
col <- ( <uint> {p.addPosNum("_col", buffer[begin:end])}
/ '\'' <singlequotedstring> '\'' {p.addPosStr("_col", buffer[begin:end])}
/ '"' <doublequotedstring> '"' {p.addPosStr("_col", buffer[begin:end])}

View file

@ -37,7 +37,6 @@ const (
rulereserved
ruleposfield
ruleuint
ruleuintrow
rulecol
rulerow
ruleopen
@ -102,7 +101,6 @@ const (
ruleAction48
ruleAction49
ruleAction50
ruleAction51
)
var rul3s = [...]string{
@ -128,7 +126,6 @@ var rul3s = [...]string{
"reserved",
"posfield",
"uint",
"uintrow",
"col",
"row",
"open",
@ -193,7 +190,6 @@ var rul3s = [...]string{
"Action48",
"Action49",
"Action50",
"Action51",
}
type token32 struct {
@ -310,7 +306,7 @@ type PQL struct {
Buffer string
buffer []rune
rules [88]func() bool
rules [86]func() bool
parse func(rule ...int) error
reset func()
Pretty bool
@ -492,20 +488,18 @@ func (p *PQL) Execute() {
case ruleAction43:
p.addPosStr("_field", buffer[begin:end])
case ruleAction44:
p.addPosNum("_row", buffer[begin:end])
case ruleAction45:
p.addPosNum("_col", buffer[begin:end])
case ruleAction45:
p.addPosStr("_col", buffer[begin:end])
case ruleAction46:
p.addPosStr("_col", buffer[begin:end])
case ruleAction47:
p.addPosStr("_col", buffer[begin:end])
case ruleAction48:
p.addPosNum("_row", buffer[begin:end])
case ruleAction48:
p.addPosStr("_row", buffer[begin:end])
case ruleAction49:
p.addPosStr("_row", buffer[begin:end])
case ruleAction50:
p.addPosStr("_row", buffer[begin:end])
case ruleAction51:
p.addPosStr("_timestamp", buffer[begin:end])
}
@ -667,7 +661,7 @@ func (p *PQL) Init() {
add(rulePegText, position13)
}
{
add(ruleAction51, position)
add(ruleAction50, position)
}
add(ruletimestamp, position12)
}
@ -753,7 +747,7 @@ func (p *PQL) Init() {
add(rulePegText, position21)
}
{
add(ruleAction48, position)
add(ruleAction47, position)
}
goto l19
l20:
@ -774,7 +768,7 @@ func (p *PQL) Init() {
}
position++
{
add(ruleAction49, position)
add(ruleAction48, position)
}
goto l19
l23:
@ -795,7 +789,7 @@ func (p *PQL) Init() {
}
position++
{
add(ruleAction50, position)
add(ruleAction49, position)
}
}
l19:
@ -2536,421 +2530,417 @@ func (p *PQL) Init() {
position, tokenIndex = position243, tokenIndex243
return false
},
/* 21 uintrow <- <(<uint> Action44)> */
nil,
/* 22 col <- <((<uint> Action45) / ('\'' <singlequotedstring> '\'' Action46) / ('"' <doublequotedstring> '"' Action47))> */
/* 21 col <- <((<uint> Action44) / ('\'' <singlequotedstring> '\'' Action45) / ('"' <doublequotedstring> '"' Action46))> */
func() bool {
position250, tokenIndex250 := position, tokenIndex
position249, tokenIndex249 := position, tokenIndex
{
position251 := position
position250 := position
{
position252, tokenIndex252 := position, tokenIndex
position251, tokenIndex251 := position, tokenIndex
{
position254 := position
position253 := position
if !_rules[ruleuint]() {
goto l253
goto l252
}
add(rulePegText, position254)
add(rulePegText, position253)
}
{
add(ruleAction45, position)
add(ruleAction44, position)
}
goto l252
l253:
position, tokenIndex = position252, tokenIndex252
goto l251
l252:
position, tokenIndex = position251, tokenIndex251
if buffer[position] != rune('\'') {
goto l256
goto l255
}
position++
{
position257 := position
position256 := position
if !_rules[rulesinglequotedstring]() {
goto l256
goto l255
}
add(rulePegText, position257)
add(rulePegText, position256)
}
if buffer[position] != rune('\'') {
goto l256
goto l255
}
position++
{
add(ruleAction45, position)
}
goto l251
l255:
position, tokenIndex = position251, tokenIndex251
if buffer[position] != rune('"') {
goto l249
}
position++
{
position258 := position
if !_rules[ruledoublequotedstring]() {
goto l249
}
add(rulePegText, position258)
}
if buffer[position] != rune('"') {
goto l249
}
position++
{
add(ruleAction46, position)
}
goto l252
l256:
position, tokenIndex = position252, tokenIndex252
if buffer[position] != rune('"') {
goto l250
}
position++
{
position259 := position
if !_rules[ruledoublequotedstring]() {
goto l250
}
add(rulePegText, position259)
}
if buffer[position] != rune('"') {
goto l250
}
position++
{
add(ruleAction47, position)
}
}
l252:
add(rulecol, position251)
l251:
add(rulecol, position250)
}
return true
l250:
position, tokenIndex = position250, tokenIndex250
l249:
position, tokenIndex = position249, tokenIndex249
return false
},
/* 23 row <- <((<uint> Action48) / ('\'' <singlequotedstring> '\'' Action49) / ('"' <doublequotedstring> '"' Action50))> */
/* 22 row <- <((<uint> Action47) / ('\'' <singlequotedstring> '\'' Action48) / ('"' <doublequotedstring> '"' Action49))> */
nil,
/* 24 open <- <('(' sp)> */
/* 23 open <- <('(' sp)> */
func() bool {
position262, tokenIndex262 := position, tokenIndex
position261, tokenIndex261 := position, tokenIndex
{
position263 := position
position262 := position
if buffer[position] != rune('(') {
goto l262
goto l261
}
position++
if !_rules[rulesp]() {
goto l262
goto l261
}
add(ruleopen, position263)
add(ruleopen, position262)
}
return true
l262:
position, tokenIndex = position262, tokenIndex262
l261:
position, tokenIndex = position261, tokenIndex261
return false
},
/* 25 close <- <(')' sp)> */
/* 24 close <- <(')' sp)> */
func() bool {
position264, tokenIndex264 := position, tokenIndex
position263, tokenIndex263 := position, tokenIndex
{
position265 := position
position264 := position
if buffer[position] != rune(')') {
goto l264
goto l263
}
position++
if !_rules[rulesp]() {
goto l264
goto l263
}
add(ruleclose, position265)
add(ruleclose, position264)
}
return true
l264:
position, tokenIndex = position264, tokenIndex264
l263:
position, tokenIndex = position263, tokenIndex263
return false
},
/* 26 sp <- <(' ' / '\t' / '\n')*> */
/* 25 sp <- <(' ' / '\t' / '\n')*> */
func() bool {
{
position267 := position
l268:
position266 := position
l267:
{
position269, tokenIndex269 := position, tokenIndex
position268, tokenIndex268 := position, tokenIndex
{
position270, tokenIndex270 := position, tokenIndex
position269, tokenIndex269 := position, tokenIndex
if buffer[position] != rune(' ') {
goto l270
}
position++
goto l269
l270:
position, tokenIndex = position269, tokenIndex269
if buffer[position] != rune('\t') {
goto l271
}
position++
goto l270
goto l269
l271:
position, tokenIndex = position270, tokenIndex270
if buffer[position] != rune('\t') {
goto l272
}
position++
goto l270
l272:
position, tokenIndex = position270, tokenIndex270
position, tokenIndex = position269, tokenIndex269
if buffer[position] != rune('\n') {
goto l269
goto l268
}
position++
}
l270:
goto l268
l269:
position, tokenIndex = position269, tokenIndex269
goto l267
l268:
position, tokenIndex = position268, tokenIndex268
}
add(rulesp, position267)
add(rulesp, position266)
}
return true
},
/* 27 comma <- <(sp ',' sp)> */
/* 26 comma <- <(sp ',' sp)> */
func() bool {
position273, tokenIndex273 := position, tokenIndex
position272, tokenIndex272 := position, tokenIndex
{
position274 := position
position273 := position
if !_rules[rulesp]() {
goto l273
goto l272
}
if buffer[position] != rune(',') {
goto l273
goto l272
}
position++
if !_rules[rulesp]() {
goto l273
goto l272
}
add(rulecomma, position274)
add(rulecomma, position273)
}
return true
l273:
position, tokenIndex = position273, tokenIndex273
l272:
position, tokenIndex = position272, tokenIndex272
return false
},
/* 28 lbrack <- <('[' sp)> */
/* 27 lbrack <- <('[' sp)> */
nil,
/* 29 rbrack <- <(sp ']' sp)> */
/* 28 rbrack <- <(sp ']' sp)> */
nil,
/* 30 IDENT <- <(([a-z] / [A-Z]) ([a-z] / [A-Z] / [0-9])*)> */
/* 29 IDENT <- <(([a-z] / [A-Z]) ([a-z] / [A-Z] / [0-9])*)> */
nil,
/* 31 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])> */
/* 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 {
position278, tokenIndex278 := position, tokenIndex
position277, tokenIndex277 := position, tokenIndex
{
position279 := position
position278 := position
if c := buffer[position]; c < rune('0') || c > rune('9') {
goto l278
goto l277
}
position++
if c := buffer[position]; c < rune('0') || c > rune('9') {
goto l278
goto l277
}
position++
if c := buffer[position]; c < rune('0') || c > rune('9') {
goto l278
goto l277
}
position++
if c := buffer[position]; c < rune('0') || c > rune('9') {
goto l278
goto l277
}
position++
if buffer[position] != rune('-') {
goto l278
goto l277
}
position++
{
position280, tokenIndex280 := position, tokenIndex
position279, tokenIndex279 := position, tokenIndex
if buffer[position] != rune('0') {
goto l281
goto l280
}
position++
goto l280
l281:
position, tokenIndex = position280, tokenIndex280
goto l279
l280:
position, tokenIndex = position279, tokenIndex279
if buffer[position] != rune('1') {
goto l278
goto l277
}
position++
}
l280:
l279:
if c := buffer[position]; c < rune('0') || c > rune('9') {
goto l278
goto l277
}
position++
if buffer[position] != rune('-') {
goto l278
goto l277
}
position++
if c := buffer[position]; c < rune('0') || c > rune('3') {
goto l278
goto l277
}
position++
if c := buffer[position]; c < rune('0') || c > rune('9') {
goto l278
goto l277
}
position++
if buffer[position] != rune('T') {
goto l278
goto l277
}
position++
if c := buffer[position]; c < rune('0') || c > rune('9') {
goto l278
goto l277
}
position++
if c := buffer[position]; c < rune('0') || c > rune('9') {
goto l278
goto l277
}
position++
if buffer[position] != rune(':') {
goto l278
goto l277
}
position++
if c := buffer[position]; c < rune('0') || c > rune('9') {
goto l278
goto l277
}
position++
if c := buffer[position]; c < rune('0') || c > rune('9') {
goto l278
goto l277
}
position++
add(ruletimestampbasicfmt, position279)
add(ruletimestampbasicfmt, position278)
}
return true
l278:
position, tokenIndex = position278, tokenIndex278
l277:
position, tokenIndex = position277, tokenIndex277
return false
},
/* 32 timestampfmt <- <(('"' timestampbasicfmt '"') / ('\'' timestampbasicfmt '\'') / timestampbasicfmt)> */
/* 31 timestampfmt <- <(('"' timestampbasicfmt '"') / ('\'' timestampbasicfmt '\'') / timestampbasicfmt)> */
func() bool {
position282, tokenIndex282 := position, tokenIndex
position281, tokenIndex281 := position, tokenIndex
{
position283 := position
position282 := position
{
position284, tokenIndex284 := position, tokenIndex
position283, tokenIndex283 := position, tokenIndex
if buffer[position] != rune('"') {
goto l284
}
position++
if !_rules[ruletimestampbasicfmt]() {
goto l284
}
if buffer[position] != rune('"') {
goto l284
}
position++
goto l283
l284:
position, tokenIndex = position283, tokenIndex283
if buffer[position] != rune('\'') {
goto l285
}
position++
if !_rules[ruletimestampbasicfmt]() {
goto l285
}
if buffer[position] != rune('"') {
if buffer[position] != rune('\'') {
goto l285
}
position++
goto l284
goto l283
l285:
position, tokenIndex = position284, tokenIndex284
if buffer[position] != rune('\'') {
goto l286
}
position++
position, tokenIndex = position283, tokenIndex283
if !_rules[ruletimestampbasicfmt]() {
goto l286
}
if buffer[position] != rune('\'') {
goto l286
}
position++
goto l284
l286:
position, tokenIndex = position284, tokenIndex284
if !_rules[ruletimestampbasicfmt]() {
goto l282
goto l281
}
}
l284:
add(ruletimestampfmt, position283)
l283:
add(ruletimestampfmt, position282)
}
return true
l282:
position, tokenIndex = position282, tokenIndex282
l281:
position, tokenIndex = position281, tokenIndex281
return false
},
/* 33 timestamp <- <(<timestampfmt> Action51)> */
/* 32 timestamp <- <(<timestampfmt> Action50)> */
nil,
/* 35 Action0 <- <{p.startCall("Set")}> */
/* 34 Action0 <- <{p.startCall("Set")}> */
nil,
/* 36 Action1 <- <{p.endCall()}> */
/* 35 Action1 <- <{p.endCall()}> */
nil,
/* 37 Action2 <- <{p.startCall("SetRowAttrs")}> */
/* 36 Action2 <- <{p.startCall("SetRowAttrs")}> */
nil,
/* 38 Action3 <- <{p.endCall()}> */
/* 37 Action3 <- <{p.endCall()}> */
nil,
/* 39 Action4 <- <{p.startCall("SetColumnAttrs")}> */
/* 38 Action4 <- <{p.startCall("SetColumnAttrs")}> */
nil,
/* 40 Action5 <- <{p.endCall()}> */
/* 39 Action5 <- <{p.endCall()}> */
nil,
/* 41 Action6 <- <{p.startCall("Clear")}> */
/* 40 Action6 <- <{p.startCall("Clear")}> */
nil,
/* 42 Action7 <- <{p.endCall()}> */
/* 41 Action7 <- <{p.endCall()}> */
nil,
/* 43 Action8 <- <{p.startCall("ClearRow")}> */
/* 42 Action8 <- <{p.startCall("ClearRow")}> */
nil,
/* 44 Action9 <- <{p.endCall()}> */
/* 43 Action9 <- <{p.endCall()}> */
nil,
/* 45 Action10 <- <{p.startCall("Store")}> */
/* 44 Action10 <- <{p.startCall("Store")}> */
nil,
/* 46 Action11 <- <{p.endCall()}> */
/* 45 Action11 <- <{p.endCall()}> */
nil,
/* 47 Action12 <- <{p.startCall("TopN")}> */
/* 46 Action12 <- <{p.startCall("TopN")}> */
nil,
/* 48 Action13 <- <{p.endCall()}> */
/* 47 Action13 <- <{p.endCall()}> */
nil,
/* 49 Action14 <- <{p.startCall("Range")}> */
/* 48 Action14 <- <{p.startCall("Range")}> */
nil,
/* 50 Action15 <- <{p.endCall()}> */
/* 49 Action15 <- <{p.endCall()}> */
nil,
nil,
/* 52 Action16 <- <{ p.startCall(buffer[begin:end] ) }> */
/* 51 Action16 <- <{ p.startCall(buffer[begin:end] ) }> */
nil,
/* 53 Action17 <- <{ p.endCall() }> */
/* 52 Action17 <- <{ p.endCall() }> */
nil,
/* 54 Action18 <- <{ p.addBTWN() }> */
/* 53 Action18 <- <{ p.addBTWN() }> */
nil,
/* 55 Action19 <- <{ p.addLTE() }> */
/* 54 Action19 <- <{ p.addLTE() }> */
nil,
/* 56 Action20 <- <{ p.addGTE() }> */
/* 55 Action20 <- <{ p.addGTE() }> */
nil,
/* 57 Action21 <- <{ p.addEQ() }> */
/* 56 Action21 <- <{ p.addEQ() }> */
nil,
/* 58 Action22 <- <{ p.addNEQ() }> */
/* 57 Action22 <- <{ p.addNEQ() }> */
nil,
/* 59 Action23 <- <{ p.addLT() }> */
/* 58 Action23 <- <{ p.addLT() }> */
nil,
/* 60 Action24 <- <{ p.addGT() }> */
/* 59 Action24 <- <{ p.addGT() }> */
nil,
/* 61 Action25 <- <{p.startConditional()}> */
/* 60 Action25 <- <{p.startConditional()}> */
nil,
/* 62 Action26 <- <{p.endConditional()}> */
/* 61 Action26 <- <{p.endConditional()}> */
nil,
/* 63 Action27 <- <{p.condAdd(buffer[begin:end])}> */
/* 62 Action27 <- <{p.condAdd(buffer[begin:end])}> */
nil,
/* 64 Action28 <- <{p.condAdd(buffer[begin:end])}> */
/* 63 Action28 <- <{p.condAdd(buffer[begin:end])}> */
nil,
/* 65 Action29 <- <{p.condAdd(buffer[begin:end])}> */
/* 64 Action29 <- <{p.condAdd(buffer[begin:end])}> */
nil,
/* 66 Action30 <- <{p.addPosStr("_start", buffer[begin:end])}> */
/* 65 Action30 <- <{p.addPosStr("_start", buffer[begin:end])}> */
nil,
/* 67 Action31 <- <{p.addPosStr("_end", buffer[begin:end])}> */
/* 66 Action31 <- <{p.addPosStr("_end", buffer[begin:end])}> */
nil,
/* 68 Action32 <- <{ p.startList() }> */
/* 67 Action32 <- <{ p.startList() }> */
nil,
/* 69 Action33 <- <{ p.endList() }> */
/* 68 Action33 <- <{ p.endList() }> */
nil,
/* 70 Action34 <- <{ p.addVal(nil) }> */
/* 69 Action34 <- <{ p.addVal(nil) }> */
nil,
/* 71 Action35 <- <{ p.addVal(true) }> */
/* 70 Action35 <- <{ p.addVal(true) }> */
nil,
/* 72 Action36 <- <{ p.addVal(false) }> */
/* 71 Action36 <- <{ p.addVal(false) }> */
nil,
/* 73 Action37 <- <{ p.addNumVal(buffer[begin:end]) }> */
/* 72 Action37 <- <{ p.addNumVal(buffer[begin:end]) }> */
nil,
/* 74 Action38 <- <{ p.addNumVal(buffer[begin:end]) }> */
/* 73 Action38 <- <{ p.addNumVal(buffer[begin:end]) }> */
nil,
/* 75 Action39 <- <{ p.addVal(buffer[begin:end]) }> */
/* 74 Action39 <- <{ p.addVal(buffer[begin:end]) }> */
nil,
/* 76 Action40 <- <{ s, _ := strconv.Unquote(buffer[begin:end]); p.addVal(s) }> */
/* 75 Action40 <- <{ s, _ := strconv.Unquote(buffer[begin:end]); p.addVal(s) }> */
nil,
/* 77 Action41 <- <{ p.addVal(buffer[begin:end]) }> */
/* 76 Action41 <- <{ p.addVal(buffer[begin:end]) }> */
nil,
/* 78 Action42 <- <{ p.addField(buffer[begin:end]) }> */
/* 77 Action42 <- <{ p.addField(buffer[begin:end]) }> */
nil,
/* 79 Action43 <- <{ p.addPosStr("_field", buffer[begin:end]) }> */
/* 78 Action43 <- <{ p.addPosStr("_field", buffer[begin:end]) }> */
nil,
/* 80 Action44 <- <{p.addPosNum("_row", buffer[begin:end])}> */
/* 79 Action44 <- <{p.addPosNum("_col", buffer[begin:end])}> */
nil,
/* 81 Action45 <- <{p.addPosNum("_col", buffer[begin:end])}> */
/* 80 Action45 <- <{p.addPosStr("_col", buffer[begin:end])}> */
nil,
/* 82 Action46 <- <{p.addPosStr("_col", buffer[begin:end])}> */
/* 81 Action46 <- <{p.addPosStr("_col", buffer[begin:end])}> */
nil,
/* 83 Action47 <- <{p.addPosStr("_col", buffer[begin:end])}> */
/* 82 Action47 <- <{p.addPosNum("_row", buffer[begin:end])}> */
nil,
/* 84 Action48 <- <{p.addPosNum("_row", buffer[begin:end])}> */
/* 83 Action48 <- <{p.addPosStr("_row", buffer[begin:end])}> */
nil,
/* 85 Action49 <- <{p.addPosStr("_row", buffer[begin:end])}> */
/* 84 Action49 <- <{p.addPosStr("_row", buffer[begin:end])}> */
nil,
/* 86 Action50 <- <{p.addPosStr("_row", buffer[begin:end])}> */
nil,
/* 87 Action51 <- <{p.addPosStr("_timestamp", buffer[begin:end])}> */
/* 85 Action50 <- <{p.addPosStr("_timestamp", buffer[begin:end])}> */
nil,
}
p.rules = _rules

View file

@ -64,6 +64,7 @@ func (sc *sliceContainers) PutContainerValues(key uint64, containerType byte, n
}
func (sc *sliceContainers) Remove(key uint64) {
statsHit("sliceContainers/Remove")
i := search64(sc.keys, key)
if i < 0 {
return
@ -73,6 +74,7 @@ func (sc *sliceContainers) Remove(key uint64) {
}
func (sc *sliceContainers) insertAt(key uint64, c *Container, i int) {
statsHit("sliceContainers/insertAt")
sc.keys = append(sc.keys, 0)
copy(sc.keys[i+1:], sc.keys[i:])
sc.keys[i] = key

View file

@ -1027,6 +1027,7 @@ func (iv interval16) runlen() int32 {
// newContainer returns a new instance of container.
func NewContainer() *Container {
statsHit("NewContainer")
return &Container{containerType: containerArray}
}
@ -1194,6 +1195,7 @@ func (c *Container) add(v uint16) (added bool) {
func (c *Container) arrayAdd(v uint16) bool {
// Optimize appending to the end of an array container.
if c.n > 0 && c.n < ArrayMaxSize && c.isArray() && c.array[c.n-1] < v {
statsHit("arrayAdd/append")
c.unmap()
c.array = append(c.array, v)
return true
@ -1207,11 +1209,13 @@ func (c *Container) arrayAdd(v uint16) bool {
// Convert to a bitmap container if too many values are in an array container.
if c.n >= ArrayMaxSize {
statsHit("arrayAdd/arrayToBitmap")
c.arrayToBitmap()
return c.bitmapAdd(v)
}
// Otherwise insert into array.
statsHit("arrayAdd/insert")
c.unmap()
i = -i - 1
c.array = append(c.array, 0)
@ -1325,6 +1329,7 @@ func (c *Container) countRuns() (r int32) {
// amount of space.
func (c *Container) optimize() {
if c.n == 0 {
statsHit("optimize/empty")
return
}
runs := c.countRuns()
@ -1341,21 +1346,33 @@ func (c *Container) optimize() {
// Then convert accordingly.
if c.isArray() {
if newType == containerBitmap {
statsHit("optimize/arrayToBitmap")
c.arrayToBitmap()
} else if newType == containerRun {
statsHit("optimize/arrayToRun")
c.arrayToRun()
} else {
statsHit("optimize/arrayUnchanged")
}
} else if c.isBitmap() {
if newType == containerArray {
statsHit("optimize/bitmapToArray")
c.bitmapToArray()
} else if newType == containerRun {
statsHit("optimize/bitmapToRun")
c.bitmapToRun()
} else {
statsHit("optimize/bitmapUnchanged")
}
} else if c.isRun() {
if newType == containerBitmap {
statsHit("optimize/runToBitmap")
c.runToBitmap()
} else if newType == containerArray {
statsHit("optimize/runToArray")
c.runToArray()
} else {
statsHit("optimize/runUnchanged")
}
}
}
@ -1425,6 +1442,7 @@ func (c *Container) bitmapRemove(v uint16) bool {
// Convert to array if we go below the threshold.
if c.n == ArrayMaxSize {
statsHit("bitmapRemove/bitmapToArray")
c.bitmapToArray()
}
return true
@ -1492,6 +1510,7 @@ func (c *Container) runMax() uint16 {
// bitmapToArray converts from bitmap format to array format.
func (c *Container) bitmapToArray() {
statsHit("bitmapToArray")
c.array = make([]uint16, 0, c.n)
c.containerType = containerArray
@ -1515,6 +1534,7 @@ func (c *Container) bitmapToArray() {
// arrayToBitmap converts from array format to bitmap format.
func (c *Container) arrayToBitmap() {
statsHit("arrayToBitmap")
c.bitmap = make([]uint64, bitmapN)
c.containerType = containerBitmap
@ -1534,6 +1554,7 @@ func (c *Container) arrayToBitmap() {
// runToBitmap converts from RLE format to bitmap format.
func (c *Container) runToBitmap() {
statsHit("runToBitmap")
c.bitmap = make([]uint64, bitmapN)
c.containerType = containerBitmap
@ -1557,6 +1578,7 @@ func (c *Container) runToBitmap() {
// bitmapToRun converts from bitmap format to RLE format.
func (c *Container) bitmapToRun() {
statsHit("bitmapToRun")
c.containerType = containerRun
// return early if empty
if c.n == 0 {
@ -1613,6 +1635,7 @@ func (c *Container) bitmapToRun() {
// arrayToRun converts from array format to RLE format.
func (c *Container) arrayToRun() {
statsHit("arrayToRun")
c.containerType = containerRun
// return early if empty
if c.n == 0 {
@ -1640,6 +1663,7 @@ func (c *Container) arrayToRun() {
// runToArray converts from RLE format to array format.
func (c *Container) runToArray() {
statsHit("runToArray")
c.containerType = containerArray
c.array = make([]uint16, 0, c.n)
@ -1661,16 +1685,20 @@ func (c *Container) runToArray() {
// Clone returns a copy of c.
func (c *Container) Clone() *Container {
statsHit("Container/Clone")
other := &Container{n: c.n, containerType: c.containerType}
switch c.containerType {
case containerArray:
statsHit("Container/Clone/Array")
other.array = make([]uint16, len(c.array))
copy(other.array, c.array)
case containerBitmap:
statsHit("Container/Clone/Bitmap")
other.bitmap = make([]uint64, len(c.bitmap))
copy(other.bitmap, c.bitmap)
case containerRun:
statsHit("Container/Clone/Run")
other.runs = make([]interval16, len(c.runs))
copy(other.runs, c.runs)
}
@ -1689,6 +1717,7 @@ func (c *Container) WriteTo(w io.Writer) (n int64, err error) {
}
func (c *Container) arrayWriteTo(w io.Writer) (n int64, err error) {
statsHit("Container/arrayWriteTo")
if len(c.array) == 0 {
return 0, nil
}
@ -1705,12 +1734,14 @@ func (c *Container) arrayWriteTo(w io.Writer) (n int64, err error) {
}
func (c *Container) bitmapWriteTo(w io.Writer) (n int64, err error) {
statsHit("Container/bitmapWriteTo")
// Write sizeof(uint64) * bitmapN bytes.
nn, err := w.Write((*[0xFFFFFFF]byte)(unsafe.Pointer(&c.bitmap[0]))[:(8 * bitmapN)])
return int64(nn), err
}
func (c *Container) runWriteTo(w io.Writer) (n int64, err error) {
statsHit("Container/runWriteTo")
if len(c.runs) == 0 {
return 0, nil
}
@ -1815,6 +1846,7 @@ func flip(a *Container) *Container { // nolint: deadcode
}
func flipArray(b *Container) *Container {
statsHit("flipArray")
// TODO: actually implement this
x := b.Clone()
x.arrayToBitmap()
@ -1822,6 +1854,7 @@ func flipArray(b *Container) *Container {
}
func flipBitmap(b *Container) *Container {
statsHit("flipBitmap")
other := &Container{bitmap: make([]uint64, bitmapN), containerType: containerBitmap}
for i, bitmap := range b.bitmap {
@ -1833,6 +1866,7 @@ func flipBitmap(b *Container) *Container {
}
func flipRun(b *Container) *Container {
statsHit("flipRun")
// TODO: actually implement this
x := b.Clone()
x.runToBitmap()
@ -1868,22 +1902,33 @@ func intersectionCount(a, b *Container) int32 {
}
func intersectionCountArrayArray(a, b *Container) (n int32) {
na, nb := len(a.array), len(b.array)
for i, j := 0, 0; i < na && j < nb; {
va, vb := a.array[i], b.array[j]
if va < vb {
i++
} else if va > vb {
statsHit("intersectionCount/ArrayArray")
ca, cb := a.array, b.array
na, nb := len(ca), len(cb)
if na == 0 || nb == 0 {
return 0
}
if na > nb {
ca, cb = cb, ca
na, nb = nb, na // nolint: ineffassign
}
j := 0
for _, va := range ca {
for cb[j] < va {
j++
} else {
if j >= nb {
return n
}
}
if cb[j] == va {
n++
i, j = i+1, j+1
}
}
return n
}
func intersectionCountArrayRun(a, b *Container) (n int32) {
statsHit("intersectionCount/ArrayRun")
na, nb := len(a.array), len(b.runs)
for i, j := 0, 0; i < na && j < nb; {
va, vb := a.array[i], b.runs[j]
@ -1900,6 +1945,7 @@ func intersectionCountArrayRun(a, b *Container) (n int32) {
}
func intersectionCountRunRun(a, b *Container) (n int32) {
statsHit("intersectionCount/RunRun")
na, nb := len(a.runs), len(b.runs)
for i, j := 0, 0; i < na && j < nb; {
va, vb := a.runs[i], b.runs[j]
@ -1931,6 +1977,7 @@ func intersectionCountRunRun(a, b *Container) (n int32) {
}
func intersectionCountBitmapRun(a, b *Container) (n int32) {
statsHit("intersectionCount/BitmapRun")
for _, iv := range b.runs {
n += a.bitmapCountRange(int32(iv.start), int32(iv.last)+1)
}
@ -1938,6 +1985,7 @@ func intersectionCountBitmapRun(a, b *Container) (n int32) {
}
func intersectionCountArrayBitmap(a, b *Container) (n int32) {
statsHit("intersectionCount/ArrayBitmap")
ln := len(b.bitmap)
for _, val := range a.array {
i := int(val >> 6)
@ -1951,6 +1999,7 @@ func intersectionCountArrayBitmap(a, b *Container) (n int32) {
}
func intersectionCountBitmapBitmap(a, b *Container) (n int32) {
statsHit("intersectionCount/BitmapBitmap")
return int32(popcountAndSlice(a.bitmap, b.bitmap))
}
@ -1983,6 +2032,7 @@ func intersect(a, b *Container) *Container {
}
func intersectArrayArray(a, b *Container) *Container {
statsHit("intersect/ArrayArray")
output := &Container{containerType: containerArray}
na, nb := len(a.array), len(b.array)
for i, j := 0, 0; i < na && j < nb; {
@ -2004,6 +2054,7 @@ func intersectArrayArray(a, b *Container) *Container {
// container. The return is always an array container (since it's guaranteed to
// be low-cardinality)
func intersectArrayRun(a, b *Container) *Container {
statsHit("intersect/ArrayRun")
output := &Container{containerType: containerArray}
na, nb := len(a.array), len(b.runs)
for i, j := 0, 0; i < na && j < nb; {
@ -2023,6 +2074,7 @@ func intersectArrayRun(a, b *Container) *Container {
// intersectRunRun computes the intersect of two run containers.
func intersectRunRun(a, b *Container) *Container {
statsHit("intersect/RunRun")
output := &Container{containerType: containerRun}
na, nb := len(a.runs), len(b.runs)
for i, j := 0, 0; i < na && j < nb; {
@ -2059,11 +2111,12 @@ func intersectRunRun(a, b *Container) *Container {
return output
}
// intersectBitmapRun returns an array container if the run container's
// cardinality is < ArrayMaxSize. Otherwise it returns a bitmap container.
// intersectBitmapRun returns an array container if either container's
// cardinality is <= ArrayMaxSize. Otherwise it returns a bitmap container.
func intersectBitmapRun(a, b *Container) *Container {
statsHit("intersect/BitmapRun")
var output *Container
if b.n < ArrayMaxSize {
if b.n <= ArrayMaxSize || a.n <= ArrayMaxSize {
// output is array container
output = &Container{containerType: containerArray}
for _, iv := range b.runs {
@ -2117,14 +2170,12 @@ func intersectBitmapRun(a, b *Container) *Container {
valast = vastart + 63
}
}
if output.n < ArrayMaxSize {
output.bitmapToArray()
}
}
return output
}
func intersectArrayBitmap(a, b *Container) *Container {
statsHit("intersect/ArrayBitmap")
output := &Container{containerType: containerArray}
for _, va := range a.array {
bmidx := va / 64
@ -2140,6 +2191,7 @@ func intersectArrayBitmap(a, b *Container) *Container {
}
func intersectBitmapBitmap(a, b *Container) *Container {
statsHit("intersect/BitmapBitmap")
// local variables added to prevent BCE checks in loop
// see https://go101.org/article/bounds-check-elimination.html
var (
@ -2191,6 +2243,7 @@ func union(a, b *Container) *Container {
}
func unionArrayArray(a, b *Container) *Container {
statsHit("union/ArrayArray")
output := &Container{containerType: containerArray}
na, nb := len(a.array), len(b.array)
for i, j := 0, 0; ; {
@ -2224,6 +2277,7 @@ func unionArrayArray(a, b *Container) *Container {
// unionArrayRun optimistically assumes that the result will be a run container,
// and converts to a bitmap or array container afterwards if necessary.
func unionArrayRun(a, b *Container) *Container {
statsHit("union/ArrayRun")
if b.n == maxContainerVal+1 {
return b.Clone()
}
@ -2281,6 +2335,7 @@ func (c *Container) runAppendInterval(v interval16) int32 {
}
func unionRunRun(a, b *Container) *Container {
statsHit("union/RunRun")
if a.n == maxContainerVal+1 {
return a.Clone()
}
@ -2315,6 +2370,7 @@ func unionRunRun(a, b *Container) *Container {
}
func unionBitmapRun(a, b *Container) *Container {
statsHit("union/BitmapRun")
if b.n == maxContainerVal+1 {
return b.Clone()
}
@ -2503,6 +2559,7 @@ func difference(a, b *Container) *Container {
// differenceArrayArray computes the difference bween two arrays.
func differenceArrayArray(a, b *Container) *Container {
statsHit("difference/ArrayArray")
output := &Container{containerType: containerArray}
na, nb := len(a.array), len(b.array)
for i, j := 0, 0; i < na; {
@ -2528,6 +2585,7 @@ func differenceArrayArray(a, b *Container) *Container {
// differenceArrayRun computes the difference of an array from a run.
func differenceArrayRun(a, b *Container) *Container {
statsHit("difference/ArrayRun")
// func (ac *arrayContainer) iandNotRun16(rc *runContainer16) container {
if a.n == 0 || b.n == 0 {
@ -2585,6 +2643,7 @@ func differenceArrayRun(a, b *Container) *Container {
// differenceBitmapRun computes the difference of an bitmap from a run.
func differenceBitmapRun(a, b *Container) *Container {
statsHit("difference/BitmapRun")
if a.n == 0 || b.n == 0 {
return a.Clone()
}
@ -2599,6 +2658,7 @@ func differenceBitmapRun(a, b *Container) *Container {
// differenceRunArray subtracts the bits in an array container from a run
// container.
func differenceRunArray(a, b *Container) *Container {
statsHit("difference/RunArray")
if a.n == 0 || b.n == 0 {
return a.Clone()
}
@ -2654,6 +2714,7 @@ RUNLOOP:
// differenceRunBitmap computes the difference of an run from a bitmap.
func differenceRunBitmap(a, b *Container) *Container {
statsHit("difference/RunBitmap")
// If a is full, difference is the flip of b.
if len(a.runs) > 0 && a.runs[0].start == 0 && a.runs[0].last == 65535 {
return flipBitmap(b)
@ -2711,6 +2772,7 @@ func differenceRunBitmap(a, b *Container) *Container {
// differenceRunRun computes the difference of two runs.
func differenceRunRun(a, b *Container) *Container {
statsHit("difference/RunRun")
if a.n == 0 || b.n == 0 {
return a.Clone()
}
@ -2774,6 +2836,7 @@ func differenceRunRun(a, b *Container) *Container {
}
func differenceArrayBitmap(a, b *Container) *Container {
statsHit("difference/ArrayBitmap")
output := &Container{containerType: containerArray}
for _, va := range a.array {
bmidx := va / 64
@ -2790,6 +2853,7 @@ func differenceArrayBitmap(a, b *Container) *Container {
}
func differenceBitmapArray(a, b *Container) *Container {
statsHit("difference/BitmapArray")
output := a.Clone()
for _, v := range b.array {
@ -2805,6 +2869,7 @@ func differenceBitmapArray(a, b *Container) *Container {
}
func differenceBitmapBitmap(a, b *Container) *Container {
statsHit("difference/BitmapBitmap")
// local variables added to prevent BCE checks in loop
// see https://go101.org/article/bounds-check-elimination.html
@ -2862,6 +2927,7 @@ func xor(a, b *Container) *Container {
}
func xorArrayArray(a, b *Container) *Container {
statsHit("xor/ArrayArray")
output := &Container{containerType: containerArray}
na, nb := len(a.array), len(b.array)
for i, j := 0, 0; i < na || j < nb; {
@ -2891,6 +2957,7 @@ func xorArrayArray(a, b *Container) *Container {
}
func xorArrayBitmap(a, b *Container) *Container {
statsHit("xor/ArrayBitmap")
output := b.Clone()
for _, v := range a.array {
if b.bitmapContains(v) {
@ -2910,6 +2977,7 @@ func xorArrayBitmap(a, b *Container) *Container {
}
func xorBitmapBitmap(a, b *Container) *Container {
statsHit("xor/BitmapBitmap")
// local variables added to prevent BCE checks in loop
// see https://go101.org/article/bounds-check-elimination.html
@ -2987,6 +3055,7 @@ func (op *op) UnmarshalBinary(data []byte) error {
if len(data) < op.size() {
return fmt.Errorf("op data out of bounds: len=%d", len(data))
}
statsHit("op/UnmarshalBinary")
// Verify checksum.
h := fnv.New32a()
@ -3011,6 +3080,7 @@ func lowbits(v uint64) uint16 { return uint16(v & 0xFFFF) }
// search32 returns the index of value in a. If value is not found, it works the
// same way as search64.
func search32(a []uint16, value uint16) int32 {
statsHit("search32")
// Optimize for elements and the last element.
n := int32(len(a))
if n == 0 {
@ -3054,6 +3124,7 @@ func search32(a []uint16, value uint16) int32 {
// since negative 0 is no different from positive 0, we offset the returned
// negative indices by 1. See the test for this function for examples.
func search64(a []uint64, value uint64) int {
statsHit("search64")
// Optimize for elements and the last element.
n := len(a)
if n == 0 {
@ -3132,6 +3203,7 @@ func (a *ErrorList) AppendWithPrefix(err error, prefix string) {
// xorArrayRun computes the exclusive or of an array and a run container.
func xorArrayRun(a, b *Container) *Container {
statsHit("xor/ArrayRun")
output := &Container{containerType: containerRun}
na, nb := len(a.array), len(b.runs)
var vb interval16
@ -3290,6 +3362,7 @@ type xorstm struct {
// xorRunRun computes the exclusive or of two run containers.
func xorRunRun(a, b *Container) *Container {
statsHit("xor/RunRun")
na, nb := len(a.runs), len(b.runs)
if na == 0 {
return b.Clone()
@ -3338,6 +3411,7 @@ func xorRunRun(a, b *Container) *Container {
// xorRunRun computes the exclusive or of a bitmap and a run container.
func xorBitmapRun(a, b *Container) *Container {
statsHit("xor/BitmapRun")
output := a.Clone()
for j := 0; j < len(b.runs); j++ {
output.bitmapXorRange(uint64(b.runs[j].start), uint64(b.runs[j].last)+1)
@ -3352,6 +3426,7 @@ func xorBitmapRun(a, b *Container) *Container {
}
func bitmapsEqual(b, c *Bitmap) error { // nolint: deadcode
statsHit("bitmapsEqual")
if b.OpWriter != c.OpWriter {
return errors.New("opWriters not equal")
}
@ -3404,6 +3479,7 @@ const (
)
func readOfficialHeader(buf []byte) (size uint32, containerTyper func(index uint, card int) byte, header, pos int, haveRuns bool, err error) {
statsHit("readOfficialHeader")
if len(buf) < 8 {
err = fmt.Errorf("buffer too small, expecting at least 8 bytes, was %d", len(buf))
return size, containerTyper, header, pos, haveRuns, err
@ -3465,6 +3541,11 @@ func readOfficialHeader(buf []byte) (size uint32, containerTyper func(index uint
// UnmarshalBinary decodes b from a binary-encoded byte slice. data can be in
// either official roaring format or Pilosa's roaring format.
func (b *Bitmap) UnmarshalBinary(data []byte) error {
if data == nil {
// Nothing to unmarshal
return nil
}
statsHit("Bitmap/UnmarshalBinary")
fileMagic := uint32(binary.LittleEndian.Uint16(data[0:2]))
if fileMagic == magicNumber { // if pilosa roaring
return errors.Wrap(b.unmarshalPilosaRoaring(data), "unmarshaling as pilosa roaring")

View file

@ -165,8 +165,8 @@ func TestRunCountRange(t *testing.T) {
}
c.add(17)
c.add(18)
c.add(19)
c.add(18)
cnt = c.runCountRange(1, 22)
if cnt != 10 {
@ -180,6 +180,11 @@ func TestRunCountRange(t *testing.T) {
if cnt != 9 {
t.Fatalf("should get 9 from multiple ranges overlapping both sides, but got: %v", cnt)
}
// verify that the disparate ops resulted in three separate runs
cnt = c.countRuns()
if cnt != 3 {
t.Fatalf("should get 3 total runs, but got: %v [%v]", cnt, c.runs)
}
}
func TestRunContains(t *testing.T) {
@ -3261,3 +3266,26 @@ func TestUnmarshalOfficialRoaring(t *testing.T) {
}
}
/*
// This function exercises an arcane edge case in dead code.
// It doesn't need to be run right now.
func TestEquals(t *testing.T) {
bma := NewBitmap()
bmr := NewBitmap()
for i := uint64(0); i < 30; i++ {
bma.Add(i)
bmr.Add(i)
}
bmr.Optimize()
bmi := bma.Intersect(bmr)
err := bitmapsEqual(bmi, bma)
if err != nil {
t.Fatalf("expected intersection to equal array")
}
err = bitmapsEqual(bmi, bmr)
if err != nil {
t.Fatalf("expected intersection to equal run")
}
}
*/

View file

@ -0,0 +1,8 @@
// +build !roaringstats
package roaring
// statsCount does nothing, because you aren't building with
// the "roaringstats" build tag.
func statsHit(string) {
}

15
roaring/roaring_stats.go Normal file
View file

@ -0,0 +1,15 @@
// +build roaringstats
package roaring
import (
"github.com/pilosa/pilosa/stats"
)
var statsEv = stats.NewExpvarStatsClient()
// statsHit increments the given stat, so we can tell how often we've hit
// that particular event.
func statsHit(name string) {
statsEv.Count(name, 1, 1)
}

View file

@ -411,13 +411,29 @@ func TestBitmap_Intersection_Empty(t *testing.T) {
}
func TestBitmap_IntersectArrayArray(t *testing.T) {
bm0 := roaring.NewFileBitmap(0, 1, 2683, 5005)
bm0 := roaring.NewFileBitmap(0, 1, 7, 9, 11, 2683, 5005)
bm1 := roaring.NewFileBitmap(0, 2683, 2684, 5000)
expected := []uint64{0, 2683}
result := bm0.Intersect(bm1)
if n := result.Count(); n != 2 {
t.Fatalf("unexpected n: %d", n)
}
for _, e := range expected {
if !result.Contains(e) {
t.Fatalf("missing value %d", e)
}
}
// confirm that it also works going the other way
result = bm1.Intersect(bm0)
if n := result.Count(); n != 2 {
t.Fatalf("unexpected n: %d", n)
}
for _, e := range expected {
if !result.Contains(e) {
t.Fatalf("missing value %d", e)
}
}
}
func TestBitmap_IntersectBitmapBitmap(t *testing.T) {
@ -689,10 +705,10 @@ func TestBitmap_Flip_After(t *testing.T) {
}
// Ensure bitmap can return the number of intersecting bits in two bitmaps.
// Ensure bitmap can return the number of intersecting bits in two arrays.
func TestBitmap_IntersectionCount_ArrayArray(t *testing.T) {
bm0 := roaring.NewFileBitmap(0, 1, 1000001, 1000002, 1000003)
bm1 := roaring.NewFileBitmap(0, 50000, 1000001, 1000002)
bm0 := roaring.NewFileBitmap(0, 1000001, 1000002, 1000003)
bm1 := roaring.NewFileBitmap(0, 50000, 999998, 999999, 1000000, 1000001, 1000002)
if n := bm0.IntersectionCount(bm1); n != 3 {
t.Fatalf("unexpected n: %d", n)
@ -1056,19 +1072,39 @@ func TestBitmapBufIterator(t *testing.T) {
}
var benchmarkBitmapIntersectionCountData struct {
a, b, r *roaring.Bitmap
// this data is used to test various operations across
// different types.
type benchmarkSampleData struct {
a1, a2, b, r1, r2 *roaring.Bitmap
}
func getBenchData() *struct{ a, b, r *roaring.Bitmap } {
data := &benchmarkBitmapIntersectionCountData
if data.a == nil {
var sampleData benchmarkSampleData
func isAllType(b *roaring.Bitmap, typ string) bool {
bi := b.Info()
for _, c := range bi.Containers {
if c.Type != typ {
return false
}
}
return true
}
func getBenchData(b *testing.B) *benchmarkSampleData {
data := &sampleData
if data.a1 == nil {
const max = (1 << 24) / 64
// Build bitmap with array container.
data.a = roaring.NewFileBitmap()
for i, n := 0, 2*roaring.ArrayMaxSize/3; i < n; i++ {
data.a.Add(uint64(rand.Intn(max)))
data.a1 = roaring.NewFileBitmap()
data.a2 = roaring.NewFileBitmap()
// two lists of different lengths
for i, n := 0, roaring.ArrayMaxSize/3; i < n; i++ {
data.a1.Add(uint64(rand.Intn(max)))
data.a2.Add(uint64(rand.Intn(max)))
}
for i, n := 0, roaring.ArrayMaxSize/3; i < n; i++ {
data.a1.Add(uint64(rand.Intn(max)))
}
// Build bitmap with bitmap container.
@ -1078,12 +1114,42 @@ func getBenchData() *struct{ a, b, r *roaring.Bitmap } {
}
// build bitmap with run container
data.r = roaring.NewFileBitmap()
data.r1 = roaring.NewFileBitmap()
for i, n := 0, MaxContainerVal; i < n; i++ {
data.r.Add(uint64(i))
data.r1.Add(uint64(i))
}
// build bitmap with multiple runs
data.r2 = roaring.NewFileBitmap()
for i, n := 0, MaxContainerVal; i < n; i++ {
data.r2.Add(uint64(i))
// break the runs up, this should produce 16 runs, which
// is small enough to make RLE tempting
if i&0xfff == 0xfff {
i += 5
}
}
data.a1.Optimize()
data.a2.Optimize()
data.b.Optimize()
data.r1.Optimize()
data.r2.Optimize()
}
if !isAllType(data.a1, "array") {
b.Fatalf("expected data.a1 to be an array, it wasn't.")
}
if !isAllType(data.a2, "array") {
b.Fatalf("expected data.a2 to be an array, it wasn't.")
}
if !isAllType(data.b, "bitmap") {
b.Fatalf("expected data.b to be a bitmap, it wasn't.")
}
if !isAllType(data.r1, "run") {
b.Fatalf("expected data.r1 to be RLE, it wasn't.")
}
if !isAllType(data.r2, "run") {
b.Fatalf("expected data.r2 to be RLE, it wasn't.")
}
return data
}
@ -1138,30 +1204,65 @@ func TestBitmap_Intersect(t *testing.T) {
}
}
func BenchmarkGetBenchData(b *testing.B) {
for i := 0; i < b.N; i++ {
sampleData = benchmarkSampleData{}
getBenchData(b)
}
}
func BenchmarkBitmap_IntersectionCount_ArrayRun(b *testing.B) {
data := getBenchData()
data := getBenchData(b)
// Reset timer & benchmark.
b.ResetTimer()
for i := 0; i < b.N; i++ {
data.a.IntersectionCount(data.r)
data.a1.IntersectionCount(data.r1)
}
}
func BenchmarkBitmap_IntersectionCount_ArrayRuns(b *testing.B) {
data := getBenchData(b)
// Reset timer & benchmark.
b.ResetTimer()
for i := 0; i < b.N; i++ {
data.a1.IntersectionCount(data.r2)
}
}
func BenchmarkBitmap_IntersectionCount_BitmapRun(b *testing.B) {
data := getBenchData()
data := getBenchData(b)
// Reset timer & benchmark.
b.ResetTimer()
for i := 0; i < b.N; i++ {
data.b.IntersectionCount(data.r)
data.b.IntersectionCount(data.r1)
}
}
func BenchmarkBitmap_IntersectionCount_BitmapRuns(b *testing.B) {
data := getBenchData(b)
// Reset timer & benchmark.
b.ResetTimer()
for i := 0; i < b.N; i++ {
data.b.IntersectionCount(data.r2)
}
}
func BenchmarkBitmap_IntersectionCount_ArrayArray(b *testing.B) {
data := getBenchData(b)
// Reset timer & benchmark.
b.ResetTimer()
for i := 0; i < b.N; i++ {
data.a1.IntersectionCount(data.a2)
data.a2.IntersectionCount(data.a1)
}
}
func BenchmarkBitmap_IntersectionCount_ArrayBitmap(b *testing.B) {
data := getBenchData()
data := getBenchData(b)
// Reset timer & benchmark.
b.ResetTimer()
for i := 0; i < b.N; i++ {
data.a.IntersectionCount(data.b)
data.a1.IntersectionCount(data.b)
}
}

View file

@ -27,7 +27,9 @@ import (
"sync"
"time"
"github.com/pilosa/pilosa/logger"
"github.com/pilosa/pilosa/roaring"
"github.com/pilosa/pilosa/stats"
"github.com/pkg/errors"
"golang.org/x/sync/errgroup"
)
@ -58,7 +60,7 @@ type Server struct { // nolint: maligned
// External
systemInfo SystemInfo
gcNotifier GCNotifier
logger Logger
logger logger.Logger
nodeID string
uri URI
@ -81,7 +83,7 @@ func (s *Server) Holder() *Holder {
// ServerOption is a functional option type for pilosa.Server
type ServerOption func(s *Server) error
func OptServerLogger(l Logger) ServerOption {
func OptServerLogger(l logger.Logger) ServerOption {
return func(s *Server) error {
s.logger = l
return nil
@ -176,7 +178,7 @@ func OptServerPrimaryTranslateStoreFunc(tf func(interface{}) TranslateStore) Ser
}
}
func OptServerStatsClient(sc StatsClient) ServerOption {
func OptServerStatsClient(sc stats.StatsClient) ServerOption {
return func(s *Server) error {
s.holder.Stats = sc
return nil
@ -258,7 +260,7 @@ func NewServer(opts ...ServerOption) (*Server, error) {
metricInterval: 0,
diagnosticInterval: 0,
logger: NopLogger,
logger: logger.NopLogger,
}
s.executor = newExecutor(optExecutorInternalQueryClient(s.defaultClient))
s.cluster.InternalClient = s.defaultClient
@ -298,6 +300,7 @@ func NewServer(opts ...ServerOption) (*Server, error) {
ID: s.nodeID,
URI: s.uri,
IsCoordinator: s.cluster.Coordinator == s.nodeID,
State: nodeStateDown,
}
s.cluster.Node = node
if s.clusterDisabled {
@ -561,7 +564,10 @@ func (s *Server) receiveMessage(m Message) error {
case *RecalculateCaches:
s.holder.recalculateCaches()
case *NodeEvent:
s.cluster.ReceiveEvent(obj)
err := s.cluster.ReceiveEvent(obj)
if err != nil {
return errors.Wrapf(err, "cluster receiving NodeEvent %v", obj)
}
case *NodeStatus:
s.handleRemoteStatus(obj)
}
@ -579,7 +585,6 @@ func (s *Server) SendSync(m Message) error {
msg = append([]byte{getMessageType(m)}, msg...)
for _, node := range s.cluster.nodes {
node := node
s.logger.Printf("SendSync to: %s", node.URI)
// Don't forward the message to ourselves.
if s.uri == node.URI {
continue
@ -600,7 +605,6 @@ func (s *Server) SendAsync(m Message) error {
// SendTo represents an implementation of Broadcaster.
func (s *Server) SendTo(to *Node, m Message) error {
s.logger.Printf("SendTo: %s", to.URI)
msg, err := s.serializer.Marshal(m)
if err != nil {
return fmt.Errorf("marshaling message: %v", err)
@ -624,7 +628,7 @@ func (s *Server) handleRemoteStatus(pb Message) {
go func() {
// Make sure the holder has opened.
<-s.holder.opened
s.holder.opened.Recv()
err := s.mergeRemoteStatus(pb.(*NodeStatus))
if err != nil {
@ -652,7 +656,7 @@ func (s *Server) mergeRemoteStatus(ns *NodeStatus) error {
// if we don't know about a field locally, log an error because
// fields should be created and synced prior to shard creation
if f == nil {
s.logger.Printf("Local Field not found: %s/%s", is.Name, fs.Name)
s.logger.Printf("local field not found: %s/%s", is.Name, fs.Name)
continue
}
if err := f.AddRemoteAvailableShards(fs.AvailableShards); err != nil {
@ -697,7 +701,7 @@ func (s *Server) monitorDiagnostics() {
s.diagnostics.CheckVersion()
err = s.diagnostics.Flush()
if err != nil {
s.logger.Printf("Diagnostics error: %s", err)
s.logger.Printf("diagnostics error: %s", err)
}
}

View file

@ -11,7 +11,7 @@
// 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 server contains the `pilosa server` subcommand which runs Pilosa
// itself. The purpose of this package is to define an easily tested Command
// object which handles interpreting configuration and setting up all the
@ -20,6 +20,7 @@
package server
import (
"bytes"
"crypto/tls"
"io"
"log"
@ -40,12 +41,14 @@ import (
"github.com/pilosa/pilosa/gopsutil"
"github.com/pilosa/pilosa/gossip"
"github.com/pilosa/pilosa/http"
"github.com/pilosa/pilosa/logger"
"github.com/pilosa/pilosa/stats"
"github.com/pilosa/pilosa/statsd"
"github.com/pkg/errors"
)
type loggerLogger interface {
pilosa.Logger
logger.Logger
Logger() *log.Logger
}
@ -139,7 +142,7 @@ func (m *Command) Start() (err error) {
go func() {
err := m.Handler.Serve()
if err != nil {
m.logger.Printf("Handler serve error: %v", err)
m.logger.Printf("handler serve error: %v", err)
}
}()
@ -148,7 +151,7 @@ func (m *Command) Start() (err error) {
return errors.Wrap(err, "opening server")
}
m.logger.Printf("Listening as %s\n", m.API.Node().URI)
m.logger.Printf("listening as %s\n", m.API.Node().URI)
return nil
}
@ -160,33 +163,37 @@ func (m *Command) Wait() error {
signal.Notify(c, os.Interrupt, syscall.SIGTERM)
select {
case sig := <-c:
m.logger.Printf("Received %s; gracefully shutting down...\n", sig.String())
m.logger.Printf("received signal '%s', gracefully shutting down...\n", sig.String())
// Second signal causes a hard shutdown.
go func() { <-c; os.Exit(1) }()
return errors.Wrap(m.Close(), "closing command")
case <-m.done:
m.logger.Printf("Server closed externally")
m.logger.Printf("server closed externally")
return nil
}
}
// setupLogger sets up the logger based on the configuration.
func (m *Command) setupLogger() error {
var err error
if m.Config.LogPath == "" {
m.logOutput = m.Stderr
} else {
m.logOutput, err = os.OpenFile(m.Config.LogPath, os.O_RDWR|os.O_CREATE|os.O_APPEND, 0600)
f, err := os.OpenFile(m.Config.LogPath, os.O_RDWR|os.O_CREATE|os.O_APPEND, 0600)
if err != nil {
return errors.Wrap(err, "opening file")
}
m.logOutput = f
err = syscall.Dup2(int(f.Fd()), int(os.Stderr.Fd()))
if err != nil {
return errors.Wrap(err, "dup2ing stderr onto logfile")
}
}
if m.Config.Verbose {
m.logger = pilosa.NewVerboseLogger(m.logOutput)
m.logger = logger.NewVerboseLogger(m.logOutput)
} else {
m.logger = pilosa.NewStandardLogger(m.logOutput)
m.logger = logger.NewStandardLogger(m.logOutput)
}
return nil
}
@ -336,7 +343,7 @@ func (m *Command) setupNetworking() error {
gossipMemberSet, err := gossip.NewMemberSet(
m.Config.Gossip,
m.API,
gossip.WithLogger(m.logger.Logger()),
gossip.WithLogOutput(&filteredWriter{logOutput: m.logOutput, v: m.Config.Verbose}),
gossip.WithTransport(m.gossipTransport),
)
if err != nil {
@ -374,14 +381,14 @@ func (m *Command) Close() error {
}
// newStatsClient creates a stats client from the config
func newStatsClient(name string, host string) (pilosa.StatsClient, error) {
func newStatsClient(name string, host string) (stats.StatsClient, error) {
switch name {
case "expvar":
return pilosa.NewExpvarStatsClient(), nil
return stats.NewExpvarStatsClient(), nil
case "statsd":
return statsd.NewStatsClient(host)
case "nop", "none":
return pilosa.NopStatsClient, nil
return stats.NopStatsClient, nil
default:
return nil, errors.Errorf("'%v' not a valid stats client, choose from [expvar, statsd, none].", name)
}
@ -407,3 +414,24 @@ func getListener(uri pilosa.URI, tlsconf *tls.Config) (ln net.Listener, err erro
return ln, nil
}
type filteredWriter struct {
v bool
logOutput io.Writer
}
// Write forwards the write to logOutput if verbose is true, or it doesn't
// contain [DEBUG] or [INFO]. This implementation isn't technically correct
// since Write could be called with only part of a log line, but I don't think
// that actually happens, so until it becomes a problem, I don't think it's
// worth dealing with the extra complexity. (jaffee)
func (f *filteredWriter) Write(p []byte) (n int, err error) {
if bytes.Contains(p, []byte("[DEBUG]")) || bytes.Contains(p, []byte("[INFO]")) {
if f.v {
return f.logOutput.Write(p)
}
} else {
return f.logOutput.Write(p)
}
return len(p), nil
}

View file

@ -12,7 +12,7 @@
// See the License for the specific language governing permissions and
// limitations under the License.
package pilosa
package stats
import (
"expvar"
@ -20,6 +20,8 @@ import (
"strings"
"sync"
"time"
"github.com/pilosa/pilosa/logger"
)
// Expvar global expvar map.
@ -52,7 +54,7 @@ type StatsClient interface {
Timing(name string, value time.Duration, rate float64)
// SetLogger Set the logger output type
SetLogger(logger Logger)
SetLogger(logger logger.Logger)
// Starts the service
Open()
@ -74,7 +76,7 @@ func (c *nopStatsClient) Gauge(name string, value float64, rate float64)
func (c *nopStatsClient) Histogram(name string, value float64, rate float64) {}
func (c *nopStatsClient) Set(name string, value string, rate float64) {}
func (c *nopStatsClient) Timing(name string, value time.Duration, rate float64) {}
func (c *nopStatsClient) SetLogger(logger Logger) {}
func (c *nopStatsClient) SetLogger(logger logger.Logger) {}
func (c *nopStatsClient) Open() {}
func (c *nopStatsClient) Close() error { return nil }
@ -149,7 +151,7 @@ func (c *expvarStatsClient) Timing(name string, value time.Duration, rate float6
}
// SetLogger has no logger.
func (c *expvarStatsClient) SetLogger(logger Logger) {
func (c *expvarStatsClient) SetLogger(logger logger.Logger) {
}
// Open no-op.
@ -221,7 +223,7 @@ func (a MultiStatsClient) Timing(name string, value time.Duration, rate float64)
}
// SetLogger Sets the StatsD logger output type.
func (a MultiStatsClient) SetLogger(logger Logger) {
func (a MultiStatsClient) SetLogger(logger logger.Logger) {
for _, c := range a {
c.SetLogger(logger)
}

View file

@ -12,7 +12,7 @@
// See the License for the specific language governing permissions and
// limitations under the License.
package pilosa_test
package stats_test
import (
"context"
@ -23,6 +23,8 @@ import (
"github.com/pilosa/pilosa"
"github.com/pilosa/pilosa/http"
"github.com/pilosa/pilosa/logger"
"github.com/pilosa/pilosa/stats"
"github.com/pilosa/pilosa/test"
)
@ -32,51 +34,51 @@ func TestMultiStatClient_Expvar(t *testing.T) {
hldr := test.MustOpenHolder()
defer hldr.Close()
c := pilosa.NewExpvarStatsClient()
ms := make(pilosa.MultiStatsClient, 1)
c := stats.NewExpvarStatsClient()
ms := make(stats.MultiStatsClient, 1)
ms[0] = c
hldr.Stats = ms
hldr.SetBit("d", "f", 0, 0)
hldr.SetBit("d", "f", 0, 1)
hldr.SetBit("d", "f", 0, ShardWidth)
hldr.SetBit("d", "f", 0, ShardWidth+2)
hldr.SetBit("d", "f", 0, pilosa.ShardWidth)
hldr.SetBit("d", "f", 0, pilosa.ShardWidth+2)
hldr.ClearBit("d", "f", 0, 1)
if pilosa.Expvar.String() != `{"index:d": {"field:f": {"view:standard": {"shard:0": {"clearBit": 1, "rows": 0, "setBit": 2}, "shard:1": {"rows": 0, "setBit": 2}}}}}` {
t.Fatalf("unexpected expvar : %s", pilosa.Expvar.String())
if stats.Expvar.String() != `{"index:d": {"field:f": {"view:standard": {"shard:0": {"clearBit": 1, "rows": 0, "setBit": 2}, "shard:1": {"rows": 0, "setBit": 2}}}}}` {
t.Fatalf("unexpected expvar : %s", stats.Expvar.String())
}
hldr.Stats.CountWithCustomTags("cc", 1, 1.0, []string{"foo:bar"})
if pilosa.Expvar.String() != `{"cc": 1, "index:d": {"field:f": {"view:standard": {"shard:0": {"clearBit": 1, "rows": 0, "setBit": 2}, "shard:1": {"rows": 0, "setBit": 2}}}}}` {
t.Fatalf("unexpected expvar : %s", pilosa.Expvar.String())
if stats.Expvar.String() != `{"cc": 1, "index:d": {"field:f": {"view:standard": {"shard:0": {"clearBit": 1, "rows": 0, "setBit": 2}, "shard:1": {"rows": 0, "setBit": 2}}}}}` {
t.Fatalf("unexpected expvar : %s", stats.Expvar.String())
}
// Gauge creates a unique key, subsequent Gauge calls will overwrite
hldr.Stats.Gauge("g", 5, 1.0)
hldr.Stats.Gauge("g", 8, 1.0)
if pilosa.Expvar.String() != `{"cc": 1, "g": 8, "index:d": {"field:f": {"view:standard": {"shard:0": {"clearBit": 1, "rows": 0, "setBit": 2}, "shard:1": {"rows": 0, "setBit": 2}}}}}` {
t.Fatalf("unexpected expvar : %s", pilosa.Expvar.String())
if stats.Expvar.String() != `{"cc": 1, "g": 8, "index:d": {"field:f": {"view:standard": {"shard:0": {"clearBit": 1, "rows": 0, "setBit": 2}, "shard:1": {"rows": 0, "setBit": 2}}}}}` {
t.Fatalf("unexpected expvar : %s", stats.Expvar.String())
}
// Set creates a unique key, subsequent sets will overwrite
hldr.Stats.Set("s", "4", 1.0)
hldr.Stats.Set("s", "7", 1.0)
if pilosa.Expvar.String() != `{"cc": 1, "g": 8, "index:d": {"field:f": {"view:standard": {"shard:0": {"clearBit": 1, "rows": 0, "setBit": 2}, "shard:1": {"rows": 0, "setBit": 2}}}}, "s": "7"}` {
t.Fatalf("unexpected expvar : %s", pilosa.Expvar.String())
if stats.Expvar.String() != `{"cc": 1, "g": 8, "index:d": {"field:f": {"view:standard": {"shard:0": {"clearBit": 1, "rows": 0, "setBit": 2}, "shard:1": {"rows": 0, "setBit": 2}}}}, "s": "7"}` {
t.Fatalf("unexpected expvar : %s", stats.Expvar.String())
}
// Record timing duration and a uniquely Set key/value
dur, _ := time.ParseDuration("123us")
hldr.Stats.Timing("tt", dur, 1.0)
if pilosa.Expvar.String() != `{"cc": 1, "g": 8, "index:d": {"field:f": {"view:standard": {"shard:0": {"clearBit": 1, "rows": 0, "setBit": 2}, "shard:1": {"rows": 0, "setBit": 2}}}}, "s": "7", "tt": 123µs}` {
t.Fatalf("unexpected expvar : %s", pilosa.Expvar.String())
if stats.Expvar.String() != `{"cc": 1, "g": 8, "index:d": {"field:f": {"view:standard": {"shard:0": {"clearBit": 1, "rows": 0, "setBit": 2}, "shard:1": {"rows": 0, "setBit": 2}}}}, "s": "7", "tt": 123µs}` {
t.Fatalf("unexpected expvar : %s", stats.Expvar.String())
}
// Expvar histogram is implemented as a gauge
hldr.Stats.Histogram("hh", 3, 1.0)
if pilosa.Expvar.String() != `{"cc": 1, "g": 8, "hh": 3, "index:d": {"field:f": {"view:standard": {"shard:0": {"clearBit": 1, "rows": 0, "setBit": 2}, "shard:1": {"rows": 0, "setBit": 2}}}}, "s": "7", "tt": 123µs}` {
t.Fatalf("unexpected expvar : %s", pilosa.Expvar.String())
if stats.Expvar.String() != `{"cc": 1, "g": 8, "hh": 3, "index:d": {"field:f": {"view:standard": {"shard:0": {"clearBit": 1, "rows": 0, "setBit": 2}, "shard:1": {"rows": 0, "setBit": 2}}}}, "s": "7", "tt": 123µs}` {
t.Fatalf("unexpected expvar : %s", stats.Expvar.String())
}
// Expvar should ignore earlier set tags from setbit
@ -92,8 +94,8 @@ func TestStatsCount_TopN(t *testing.T) {
hldr.SetBit("d", "f", 0, 0)
hldr.SetBit("d", "f", 0, 1)
hldr.SetBit("d", "f", 0, ShardWidth)
hldr.SetBit("d", "f", 0, ShardWidth+2)
hldr.SetBit("d", "f", 0, pilosa.ShardWidth)
hldr.SetBit("d", "f", 0, pilosa.ShardWidth+2)
// Execute query.
called := false
@ -311,11 +313,11 @@ func (s *MockStats) CountWithCustomTags(name string, value int64, rate float64,
}
func (c *MockStats) Tags() []string { return nil }
func (c *MockStats) WithTags(tags ...string) pilosa.StatsClient { return c }
func (c *MockStats) WithTags(tags ...string) stats.StatsClient { return c }
func (c *MockStats) Gauge(name string, value float64, rate float64) {}
func (c *MockStats) Histogram(name string, value float64, rate float64) {}
func (c *MockStats) Set(name string, value string, rate float64) {}
func (c *MockStats) Timing(name string, value time.Duration, rate float64) {}
func (c *MockStats) SetLogger(logger pilosa.Logger) {}
func (c *MockStats) SetLogger(logger logger.Logger) {}
func (c *MockStats) Open() {}
func (c *MockStats) Close() error { return nil }

View file

@ -19,7 +19,8 @@ import (
"time"
"github.com/DataDog/datadog-go/statsd"
"github.com/pilosa/pilosa"
"github.com/pilosa/pilosa/logger"
"github.com/pilosa/pilosa/stats"
)
// StatsD protocol wrapper using the DataDog library that added Tags to the StatsD protocol
@ -34,13 +35,13 @@ const (
)
// Ensure client implements interface.
var _ pilosa.StatsClient = &statsClient{}
var _ stats.StatsClient = &statsClient{}
// statsClient represents a StatsD implementation of pilosa.statsClient.
type statsClient struct {
client *statsd.Client
tags []string
logger pilosa.Logger
logger logger.Logger
}
// NewStatsClient returns a new instance of StatsClient.
@ -52,7 +53,7 @@ func NewStatsClient(host string) (*statsClient, error) {
return &statsClient{
client: c,
logger: pilosa.NopLogger,
logger: logger.NopLogger,
}, nil
}
@ -70,7 +71,7 @@ func (c *statsClient) Tags() []string {
}
// WithTags returns a new client with additional tags appended.
func (c *statsClient) WithTags(tags ...string) pilosa.StatsClient {
func (c *statsClient) WithTags(tags ...string) stats.StatsClient {
return &statsClient{
client: c.client,
tags: unionStringSlice(c.tags, tags),
@ -122,7 +123,7 @@ func (c *statsClient) Timing(name string, value time.Duration, rate float64) {
}
// SetLogger sets the logger for client.
func (c *statsClient) SetLogger(logger pilosa.Logger) {
func (c *statsClient) SetLogger(logger logger.Logger) {
c.logger = logger
}

View file

@ -15,6 +15,7 @@ import (
"time"
"github.com/cespare/xxhash"
"github.com/pilosa/pilosa/logger"
"github.com/pkg/errors"
)
@ -68,7 +69,7 @@ type TranslateFile struct {
Path string
mapSize int
logger Logger
logger logger.Logger
// If non-nil, data is streamed from a primary and this is a read-only store.
PrimaryTranslateStore TranslateStore
primaryID string // unique ID used to identify the primary store
@ -89,7 +90,7 @@ func OptTranslateFileMapSize(mapSize int) TranslateFileOption {
return nil
}
}
func OptTranslateFileLogger(l Logger) TranslateFileOption {
func OptTranslateFileLogger(l logger.Logger) TranslateFileOption {
return func(s *TranslateFile) error {
s.logger = l
return nil
@ -116,7 +117,7 @@ func NewTranslateFile(opts ...TranslateFileOption) *TranslateFile {
mapSize: defaultMapSize,
logger: NopLogger,
logger: logger.NopLogger,
replicationClosing: make(chan struct{}),
primaryStoreEvents: make(chan primaryStoreEvent),
@ -194,12 +195,11 @@ func (s *TranslateFile) handlePrimaryStoreEvent(ev primaryStoreEvent) error {
}
// Stop translate store replication.
s.logger.Printf("stop monitor replication")
close(s.replicationClosing)
s.repWG.Wait()
// Set the primary node for translate store replication.
s.logger.Printf("set primary translate store to %s", ev.id)
s.logger.Debugf("set primary translate store to %s", ev.id)
s.primaryID = ev.id
if ev.id == "" {
s.PrimaryTranslateStore = nil
@ -208,7 +208,6 @@ func (s *TranslateFile) handlePrimaryStoreEvent(ev primaryStoreEvent) error {
}
// Start translate store replication. Stream from primary, if available.
s.logger.Printf("start monitor replication")
if s.PrimaryTranslateStore != nil {
s.replicationClosing = make(chan struct{})
s.repWG.Add(1)
@ -385,7 +384,6 @@ func (s *TranslateFile) monitorReplication() {
// monitorPrimaryStoreEvents is executed in a separate goroutine and listens for changes
// to the primary store assignment.
func (s *TranslateFile) monitorPrimaryStoreEvents() {
s.logger.Printf("monitor primary store events")
// Keep handling events until the store closes.
for {
select {
@ -403,7 +401,7 @@ func (s *TranslateFile) replicate(ctx context.Context) error {
off := s.size()
// Connect to remote primary.
s.logger.Printf("pilosa: replicating from offset %d", off)
s.logger.Debugf("pilosa: replicating from offset %d", off)
rc, err := s.PrimaryTranslateStore.Reader(ctx, off)
if err != nil {
return err

View file

@ -809,7 +809,7 @@ func NewTranslateFile() *TranslateFile {
}
f.Close()
s := &TranslateFile{TranslateFile: pilosa.NewTranslateFile(pilosa.OptTranslateFileMapSize(2 << 25))}
s := &TranslateFile{TranslateFile: pilosa.NewTranslateFile(pilosa.OptTranslateFileMapSize(2 << 26))}
s.Path = f.Name()
return s
}

10
view.go
View file

@ -22,8 +22,10 @@ import (
"strings"
"sync"
"github.com/pilosa/pilosa/logger"
"github.com/pilosa/pilosa/pql"
"github.com/pilosa/pilosa/roaring"
"github.com/pilosa/pilosa/stats"
"github.com/pkg/errors"
)
@ -50,9 +52,9 @@ type view struct {
fragments map[uint64]*fragment
broadcaster broadcaster
stats StatsClient
stats stats.StatsClient
rowAttrStore AttrStore
logger Logger
logger logger.Logger
}
// newView returns a new instance of View.
@ -70,8 +72,8 @@ func newView(path, index, field, name string, fieldOptions FieldOptions) *view {
fragments: make(map[uint64]*fragment),
broadcaster: NopBroadcaster,
stats: NopStatsClient,
logger: NopLogger,
stats: stats.NopStatsClient,
logger: logger.NopLogger,
}
}