mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-09-05 08:10:50 +00:00
commit
f9513ed0bc
14 changed files with 885 additions and 20 deletions
|
|
@ -67,7 +67,7 @@ jobs:
|
|||
name: golang
|
||||
steps:
|
||||
- run: '[[ -n $CIRCLE_PULL_REQUEST ]] || circleci step halt || true' # Skip if this is not a pull request
|
||||
- run: curl https://moleculacorp:$GITHUB_PERSONAL_ACCESS_TOKEN@api.github.com/repos/molecula/pilosa/pulls/$(basename $CIRCLE_PULL_REQUEST) | jq "[.labels[] | .name | startswith(\"changelog\")] | any" -e
|
||||
- run: curl https://moleculacorp:$GITHUB_PERSONAL_ACCESS_TOKEN@api.github.com/repos/molecula/featurebase/pulls/$(basename $CIRCLE_PULL_REQUEST) | jq "[.labels[] | .name | startswith(\"changelog\")] | any" -e
|
||||
test-build-arm:
|
||||
executor:
|
||||
name: golang
|
||||
|
|
|
|||
259
api.go
259
api.go
|
|
@ -2179,6 +2179,263 @@ func (api *API) TranslateIndexDB(ctx context.Context, indexName string, partitio
|
|||
_, err := store.ReadFrom(rd)
|
||||
return err
|
||||
}
|
||||
func (api *API) mutexCheckThisNode(ctx context.Context, qcx *Qcx, indexName string, fieldName string) (map[uint64]map[uint64][]uint64, error) {
|
||||
index := api.holder.Index(indexName)
|
||||
if index == nil {
|
||||
return nil, newNotFoundError(ErrIndexNotFound, indexName)
|
||||
}
|
||||
field := index.Field(fieldName)
|
||||
if field == nil {
|
||||
return nil, newNotFoundError(ErrFieldNotFound, fieldName)
|
||||
}
|
||||
return field.MutexCheck(ctx, qcx)
|
||||
}
|
||||
|
||||
// mergeIDLists merges a list of numeric IDs into another list, removing
|
||||
// duplicates.
|
||||
func mergeIDLists(dst []uint64, src []uint64) []uint64 {
|
||||
dst = append(dst, src...)
|
||||
sort.Slice(dst, func(i, j int) bool {
|
||||
return dst[i] < dst[j]
|
||||
})
|
||||
// dedup.
|
||||
n := 0
|
||||
prev := dst[0]
|
||||
for i := 0; i < len(dst); i++ {
|
||||
if dst[i] != prev {
|
||||
dst[n] = dst[i]
|
||||
n++
|
||||
}
|
||||
prev = dst[i]
|
||||
}
|
||||
return dst
|
||||
}
|
||||
|
||||
// mergeKeyLists merges a list of string IDs into another list, removing
|
||||
// duplicates.
|
||||
func mergeKeyLists(dst []string, src []string) []string {
|
||||
dst = append(dst, src...)
|
||||
sort.Slice(dst, func(i, j int) bool {
|
||||
return dst[i] < dst[j]
|
||||
})
|
||||
// dedup.
|
||||
n := 0
|
||||
prev := dst[0]
|
||||
for i := 0; i < len(dst); i++ {
|
||||
if dst[i] != prev {
|
||||
dst[n] = dst[i]
|
||||
n++
|
||||
}
|
||||
prev = dst[i]
|
||||
}
|
||||
return dst
|
||||
}
|
||||
|
||||
// MutexCheckNode checks for collisions in a given mutex field. The response is
|
||||
// a map[shard]map[column]values, not translated.
|
||||
func (api *API) MutexCheckNode(ctx context.Context, qcx *Qcx, indexName string, fieldName string) (map[uint64]map[uint64][]uint64, error) {
|
||||
if err := api.validate(apiMutexCheck); err != nil {
|
||||
return nil, errors.Wrap(err, "validating api method")
|
||||
}
|
||||
return api.mutexCheckThisNode(ctx, qcx, indexName, fieldName)
|
||||
}
|
||||
|
||||
// MutexCheck checks a named field for mutex violations, returning a
|
||||
// map of record IDs to values for records that have multiple values in the
|
||||
// field. The return will be one of:
|
||||
// map[uint64][]uint64 // unkeyed index, unkeyed field
|
||||
// map[uint64][]string // unkeyed index, keyed field
|
||||
// map[string][]uint64 // keyed index, unkeyed field
|
||||
// map[string][]string // keyed index, keyed field
|
||||
func (api *API) MutexCheck(ctx context.Context, qcx *Qcx, indexName string, fieldName string) (result interface{}, err error) {
|
||||
if err = api.validate(apiMutexCheck); err != nil {
|
||||
return nil, errors.Wrap(err, "validating api method")
|
||||
}
|
||||
index, err := api.Index(ctx, indexName)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
field, err := api.Field(ctx, indexName, fieldName)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if field.Type() != FieldTypeMutex {
|
||||
return nil, errors.New("can only check mutex state for mutex fields")
|
||||
}
|
||||
nodes := Nodes(api.cluster.nodes).Clone()
|
||||
eg, _ := errgroup.WithContext(ctx)
|
||||
|
||||
results := make([]map[uint64]map[uint64][]uint64, len(nodes))
|
||||
myID := api.Node().ID
|
||||
// vprint.VV("MyID %#v\n", myID)
|
||||
for i, node := range nodes {
|
||||
i := i // loop variable shadowing is a war crime
|
||||
// vprint.VV("Compare %#v with %v", node.ID, myID)
|
||||
if node.ID != myID {
|
||||
node := node // loop variable shadowing again
|
||||
eg.Go(func() (err error) {
|
||||
results[i], err = api.server.defaultClient.MutexCheck(ctx, &node.URI, indexName, fieldName)
|
||||
return err
|
||||
})
|
||||
} else {
|
||||
eg.Go(func() (err error) {
|
||||
results[i], err = api.mutexCheckThisNode(ctx, qcx, indexName, fieldName)
|
||||
return err
|
||||
})
|
||||
}
|
||||
}
|
||||
err = eg.Wait()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
// We now have a series of maps from shards to maps of record IDs to
|
||||
// values. But wait! Either the field, or the index, might be using keys,
|
||||
// and want those translated. So we have to translate those.
|
||||
useIndexKeys := index.Keys()
|
||||
useFieldKeys := field.Keys()
|
||||
var indexKeys = map[uint64]string{}
|
||||
var fieldKeys = map[uint64]string{}
|
||||
var indexIDs []uint64
|
||||
var fieldIDs []uint64
|
||||
// build translation tables, if we need them
|
||||
untranslated := "untranslated"
|
||||
// We don't know which of four map types we want to be working with,
|
||||
// but what we can do is make a function which works with that map type
|
||||
// given the raw integer values, and is a closure with an already-created
|
||||
// map which has already been stashed in `result`. Because maps are
|
||||
// reference-y, this should actually work.
|
||||
var process func(uint64, []uint64)
|
||||
if useIndexKeys || useFieldKeys {
|
||||
for _, nodeResults := range results {
|
||||
for _, shardResults := range nodeResults {
|
||||
for record, values := range shardResults {
|
||||
if useIndexKeys {
|
||||
if _, ok := indexKeys[record]; !ok {
|
||||
indexKeys[record] = untranslated
|
||||
indexIDs = append(indexIDs, record)
|
||||
}
|
||||
}
|
||||
if useFieldKeys {
|
||||
for _, value := range values {
|
||||
if _, ok := fieldKeys[value]; !ok {
|
||||
fieldKeys[value] = untranslated
|
||||
fieldIDs = append(fieldIDs, value)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
// we now have lists of the keys, so...
|
||||
if useIndexKeys {
|
||||
indexKeyList, err := api.cluster.translateIndexIDs(ctx, indexName, indexIDs)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "translating index keys")
|
||||
}
|
||||
if len(indexKeyList) != len(indexIDs) {
|
||||
return nil, fmt.Errorf("translating %d record IDs, got %d keys", len(indexIDs), len(indexKeyList))
|
||||
}
|
||||
for i := range indexIDs {
|
||||
indexKeys[indexIDs[i]] = indexKeyList[i]
|
||||
}
|
||||
}
|
||||
if useFieldKeys {
|
||||
fieldKeyList, err := api.cluster.translateFieldListIDs(field, fieldIDs)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "translating index keys")
|
||||
}
|
||||
if len(fieldKeyList) != len(fieldIDs) {
|
||||
return nil, fmt.Errorf("translating %d IDs, got %d keys", len(indexIDs), len(fieldKeyList))
|
||||
}
|
||||
for i := range fieldIDs {
|
||||
fieldKeys[fieldIDs[i]] = fieldKeyList[i]
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
// define the process functions. separated from above code just to make
|
||||
// it easier to follow/compare them.
|
||||
if useIndexKeys {
|
||||
if useFieldKeys {
|
||||
outMap := make(map[string][]string)
|
||||
var valueKeys []string
|
||||
result = outMap
|
||||
process = func(recordID uint64, valueIDs []uint64) {
|
||||
valueKeys = valueKeys[:0]
|
||||
for _, id := range valueIDs {
|
||||
valueKeys = append(valueKeys, fieldKeys[id])
|
||||
}
|
||||
record := indexKeys[recordID]
|
||||
if existing, ok := outMap[record]; ok {
|
||||
outMap[record] = mergeKeyLists(existing, valueKeys)
|
||||
} else {
|
||||
// The append is so we can reuse this buffer safely,
|
||||
// which matters if there's replication, because many
|
||||
// cases won't need to copy the buffer, they'll just
|
||||
// copy individual things from it.
|
||||
outMap[record] = append([]string{}, valueKeys...)
|
||||
}
|
||||
}
|
||||
} else {
|
||||
outMap := make(map[string][]uint64)
|
||||
result = outMap
|
||||
process = func(recordID uint64, values []uint64) {
|
||||
record := indexKeys[recordID]
|
||||
if existing, ok := outMap[record]; ok {
|
||||
outMap[record] = mergeIDLists(existing, values)
|
||||
} else {
|
||||
outMap[record] = values
|
||||
}
|
||||
}
|
||||
}
|
||||
} else {
|
||||
if useFieldKeys {
|
||||
outMap := make(map[uint64][]string)
|
||||
var valueKeys []string
|
||||
result = outMap
|
||||
process = func(record uint64, valueIDs []uint64) {
|
||||
valueKeys = valueKeys[:0]
|
||||
for _, id := range valueIDs {
|
||||
valueKeys = append(valueKeys, fieldKeys[id])
|
||||
}
|
||||
if existing, ok := outMap[record]; ok {
|
||||
outMap[record] = mergeKeyLists(existing, valueKeys)
|
||||
} else {
|
||||
// The append is so we can reuse this buffer safely,
|
||||
// which matters if there's replication, because many
|
||||
// cases won't need to copy the buffer, they'll just
|
||||
// copy individual things from it.
|
||||
outMap[record] = append([]string{}, valueKeys...)
|
||||
}
|
||||
}
|
||||
} else {
|
||||
outMap := make(map[uint64][]uint64)
|
||||
result = outMap
|
||||
process = func(record uint64, values []uint64) {
|
||||
if existing, ok := outMap[record]; ok {
|
||||
outMap[record] = mergeIDLists(existing, values)
|
||||
} else {
|
||||
outMap[record] = values
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
for _, nodeResults := range results {
|
||||
if len(nodeResults) == 0 {
|
||||
continue
|
||||
}
|
||||
for _, v := range nodeResults {
|
||||
if len(v) == 0 {
|
||||
continue
|
||||
}
|
||||
for record, values := range v {
|
||||
process(record, values)
|
||||
}
|
||||
}
|
||||
}
|
||||
return result, nil
|
||||
}
|
||||
|
||||
// TranslateFieldDB is an internal function to load the field keys database
|
||||
func (api *API) TranslateFieldDB(ctx context.Context, indexName, fieldName string, rd io.Reader) error {
|
||||
|
|
@ -2249,11 +2506,13 @@ const (
|
|||
apiIDReserve
|
||||
apiIDCommit
|
||||
apiIDReset
|
||||
apiMutexCheck //Backported from 1681 Caution
|
||||
)
|
||||
|
||||
var methodsCommon = map[apiMethod]struct{}{
|
||||
apiClusterMessage: {},
|
||||
apiSetCoordinator: {},
|
||||
apiMutexCheck: {},
|
||||
}
|
||||
|
||||
var methodsResizing = map[apiMethod]struct{}{
|
||||
|
|
|
|||
315
api_test.go
315
api_test.go
|
|
@ -28,6 +28,7 @@ import (
|
|||
"github.com/pilosa/pilosa/v2/boltdb"
|
||||
"github.com/pilosa/pilosa/v2/http"
|
||||
"github.com/pilosa/pilosa/v2/server"
|
||||
"github.com/pilosa/pilosa/v2/shardwidth"
|
||||
"github.com/pilosa/pilosa/v2/test"
|
||||
)
|
||||
|
||||
|
|
@ -622,3 +623,317 @@ func TestAPI_ClearFlagForImportAndImportValues(t *testing.T) {
|
|||
panic(fmt.Sprintf("expected %v, observed %v starting acct0 balance", acct0bal, 0))
|
||||
}
|
||||
}
|
||||
|
||||
type mutexCheckIndex struct {
|
||||
index *pilosa.Index
|
||||
indexName string
|
||||
createdAt int64
|
||||
fields map[bool]mutexCheckField
|
||||
}
|
||||
|
||||
type mutexCheckField struct {
|
||||
fieldName string
|
||||
field *pilosa.Field
|
||||
createdAt int64
|
||||
}
|
||||
|
||||
func TestAPI_MutexCheck(t *testing.T) {
|
||||
c := test.MustRunCluster(t, 3)
|
||||
defer c.Close()
|
||||
|
||||
m0 := c.GetNode(0)
|
||||
nodesByID := make(map[string]*test.Command, 3)
|
||||
qcxsByID := make(map[string]*pilosa.Qcx, 3)
|
||||
for i := 0; i < 3; i++ {
|
||||
node := c.GetNode(i)
|
||||
id := node.API.Node().ID
|
||||
nodesByID[id] = node
|
||||
}
|
||||
|
||||
indexes := make(map[bool]mutexCheckIndex)
|
||||
|
||||
ctx := context.Background()
|
||||
for _, keyedIndex := range []bool{false, true} {
|
||||
indexName := fmt.Sprintf("i%t", keyedIndex)
|
||||
index, err := m0.API.CreateIndex(ctx, indexName, pilosa.IndexOptions{Keys: keyedIndex, TrackExistence: true})
|
||||
if err != nil {
|
||||
t.Fatalf("creating index: %v", err)
|
||||
}
|
||||
if index.CreatedAt() == 0 {
|
||||
t.Fatal("index createdAt is empty")
|
||||
}
|
||||
indexData := mutexCheckIndex{indexName: indexName, index: index, fields: make(map[bool]mutexCheckField), createdAt: index.CreatedAt()}
|
||||
for _, keyedField := range []bool{false, true} {
|
||||
fieldName := fmt.Sprintf("f%t", keyedField)
|
||||
var field *pilosa.Field
|
||||
if keyedField {
|
||||
field, err = m0.API.CreateField(ctx, indexName, fieldName, pilosa.OptFieldTypeMutex(pilosa.CacheTypeNone, 0), pilosa.OptFieldKeys())
|
||||
} else {
|
||||
field, err = m0.API.CreateField(ctx, indexName, fieldName, pilosa.OptFieldTypeMutex(pilosa.CacheTypeNone, 0))
|
||||
}
|
||||
if err != nil {
|
||||
t.Fatalf("creating field: %v", err)
|
||||
}
|
||||
if field.CreatedAt() == 0 {
|
||||
t.Fatal("field createdAt is empty")
|
||||
}
|
||||
indexData.fields[keyedField] = mutexCheckField{fieldName: fieldName, field: field, createdAt: field.CreatedAt()}
|
||||
}
|
||||
indexes[keyedIndex] = indexData
|
||||
}
|
||||
|
||||
rowIDs := []uint64{0, 1, 2, 3}
|
||||
colIDs := []uint64{0, 1, 2, 3}
|
||||
rowKeysBase := []string{"v0", "v1", "v2", "v3"}
|
||||
colKeysBase := []string{"c0", "c1", "c2", "c3"}
|
||||
|
||||
const nShards = 10
|
||||
|
||||
// now, try the same thing for each combination of keyed/unkeyed. we
|
||||
// share code between keyed/unkeyed fields, but for indexes, the logic
|
||||
// is fundamentally different because we can't know shards in advance.
|
||||
indexData := indexes[false]
|
||||
for keyedField, fieldData := range indexData.fields {
|
||||
fmt.Println("running:", indexData.indexName, fieldData.fieldName)
|
||||
t.Run(fmt.Sprintf("%s-%s", indexData.indexName, fieldData.fieldName), func(t *testing.T) {
|
||||
for id, node := range nodesByID {
|
||||
qcxsByID[id] = node.API.Txf().NewQcx()
|
||||
}
|
||||
for shard := uint64(0); shard < nShards; shard++ {
|
||||
// restore row/col ID values which can get altered by imports
|
||||
for i := range rowIDs {
|
||||
rowIDs[i] = uint64(i)
|
||||
colIDs[i] = (shard << shardwidth.Exponent) + uint64(i) + (shard % 4)
|
||||
}
|
||||
req := &pilosa.ImportRequest{
|
||||
Index: indexData.indexName,
|
||||
IndexCreatedAt: indexData.createdAt,
|
||||
Field: fieldData.fieldName,
|
||||
FieldCreatedAt: fieldData.createdAt,
|
||||
Shard: shard,
|
||||
ColumnIDs: colIDs,
|
||||
}
|
||||
if keyedField {
|
||||
req.RowKeys = rowKeysBase
|
||||
} else {
|
||||
req.RowIDs = rowIDs
|
||||
}
|
||||
nodesForShard, err := m0.API.ShardNodes(ctx, indexData.indexName, shard)
|
||||
if err != nil {
|
||||
t.Fatalf("obtaining shard list: %v", err)
|
||||
}
|
||||
if len(nodesForShard) < 1 {
|
||||
t.Fatalf("no nodes for shard %d", shard)
|
||||
}
|
||||
node := nodesByID[nodesForShard[0].ID]
|
||||
if err := node.API.Import(ctx, qcxsByID[nodesForShard[0].ID], req); err != nil {
|
||||
t.Fatalf("importing data: %v", err)
|
||||
}
|
||||
}
|
||||
// and then we break the mutex and close the Qcxs
|
||||
for id, node := range nodesByID {
|
||||
field, err := node.API.Field(ctx, indexData.indexName, fieldData.fieldName)
|
||||
if err != nil {
|
||||
t.Fatalf("requesting field %s from node %s: %v", fieldData.fieldName, id, err)
|
||||
}
|
||||
pilosa.CorruptAMutex(t, field, qcxsByID[id])
|
||||
err = qcxsByID[id].Finish()
|
||||
if err != nil {
|
||||
t.Fatalf("closing out transaction on node %s: %v", id, err)
|
||||
}
|
||||
}
|
||||
qcx := m0.API.Txf().NewQcx()
|
||||
defer qcx.Abort()
|
||||
|
||||
results, err := m0.API.MutexCheck(ctx, qcx, indexData.indexName, fieldData.fieldName)
|
||||
if err != nil {
|
||||
t.Fatalf("checking mutexes: %v", err)
|
||||
}
|
||||
// first two shards of each group of 4 should have a collision in
|
||||
// position 1
|
||||
expected := map[uint64]bool{
|
||||
(0 << shardwidth.Exponent) + 1: true,
|
||||
(1 << shardwidth.Exponent) + 1: true,
|
||||
(4 << shardwidth.Exponent) + 1: true,
|
||||
(5 << shardwidth.Exponent) + 1: true,
|
||||
(8 << shardwidth.Exponent) + 1: true,
|
||||
(9 << shardwidth.Exponent) + 1: true,
|
||||
}
|
||||
if keyedField {
|
||||
mapped, ok := results.(map[uint64][]string)
|
||||
if !ok {
|
||||
t.Fatalf("expected map[uint64][]string, got %T", results)
|
||||
}
|
||||
seen := 0
|
||||
for k, v := range mapped {
|
||||
seen++
|
||||
if !expected[k] {
|
||||
t.Fatalf("expected all collisions to be 1 shards (s %% 4 in [0,1]), got %d", k)
|
||||
}
|
||||
if len(v) != 2 {
|
||||
t.Fatalf("expected exactly two collisions")
|
||||
}
|
||||
}
|
||||
if seen != len(expected) {
|
||||
t.Fatalf("expected exactly %d records to have collisions", len(expected))
|
||||
}
|
||||
} else {
|
||||
mapped, ok := results.(map[uint64][]uint64)
|
||||
if !ok {
|
||||
t.Fatalf("expected map[uint64][]uint64, got %T", results)
|
||||
}
|
||||
seen := 0
|
||||
for k, v := range mapped {
|
||||
seen++
|
||||
if !expected[k] {
|
||||
t.Fatalf("expected all collisions to be 1 shards (s %% 4 in [0,1]), got %d", k)
|
||||
}
|
||||
if len(v) != 2 {
|
||||
t.Fatalf("expected exactly two collisions")
|
||||
}
|
||||
}
|
||||
if seen != len(expected) {
|
||||
t.Fatalf("expected exactly %d records to have collisions", len(expected))
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
indexData = indexes[true]
|
||||
for keyedField, fieldData := range indexData.fields {
|
||||
t.Run(fmt.Sprintf("%s-%s", indexData.indexName, fieldData.fieldName), func(t *testing.T) {
|
||||
for id, node := range nodesByID {
|
||||
qcxsByID[id] = node.API.Txf().NewQcx()
|
||||
}
|
||||
req := &pilosa.ImportRequest{
|
||||
Index: indexData.indexName,
|
||||
IndexCreatedAt: indexData.createdAt,
|
||||
Field: fieldData.fieldName,
|
||||
FieldCreatedAt: fieldData.createdAt,
|
||||
Shard: 0, // ignored when using keys
|
||||
}
|
||||
rowKeys := make([]string, 0, len(rowKeysBase)*nShards)
|
||||
colKeys := make([]string, 0, len(rowKeysBase)*nShards)
|
||||
rowIDs = rowIDs[:0]
|
||||
for shard := uint64(0); shard < nShards; shard++ {
|
||||
for i := range rowKeysBase {
|
||||
colKeys = append(colKeys, fmt.Sprintf("s%d-%s", shard, colKeysBase[i]))
|
||||
if keyedField {
|
||||
rowKeys = append(rowKeys, rowKeysBase[i])
|
||||
} else {
|
||||
rowIDs = append(rowIDs, uint64(i))
|
||||
}
|
||||
}
|
||||
}
|
||||
req.ColumnKeys = colKeys
|
||||
if keyedField {
|
||||
req.RowKeys = rowKeys
|
||||
} else {
|
||||
req.RowIDs = rowIDs
|
||||
}
|
||||
var id string
|
||||
var node *test.Command
|
||||
for id, node = range nodesByID {
|
||||
break
|
||||
}
|
||||
if err := node.API.Import(ctx, qcxsByID[id], req); err != nil {
|
||||
t.Fatalf("importing data: %v", err)
|
||||
}
|
||||
expected, err := node.API.FindIndexKeys(ctx, indexData.indexName, colKeys...)
|
||||
if err != nil {
|
||||
t.Fatalf("looking up index keys: %v", err)
|
||||
}
|
||||
for key, id := range expected {
|
||||
// CorruptAMutex should only corrupt things in position 1 of their
|
||||
// shards...
|
||||
if id%(1<<shardwidth.Exponent) != 1 {
|
||||
delete(expected, key)
|
||||
}
|
||||
}
|
||||
if keyedField {
|
||||
fieldValues, err := node.API.FindFieldKeys(ctx, indexData.indexName, fieldData.fieldName, rowKeys...)
|
||||
if err != nil {
|
||||
t.Fatalf("looking up field keys: %v", err)
|
||||
}
|
||||
// Figure out which key got the value 3, delete any records
|
||||
// which would have had that key, because they won't be
|
||||
// conflicts.
|
||||
for key, value := range fieldValues {
|
||||
if value == 3 {
|
||||
for offset, baseKey := range rowKeysBase {
|
||||
if baseKey == key {
|
||||
for i := offset; i < len(rowKeys); i += len(rowKeysBase) {
|
||||
delete(expected, colKeys[i])
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
} else {
|
||||
// we set rowKeys to 0-1-2-... for rowKeysBase items, which
|
||||
// tells us which keys we expect to be 3 already.
|
||||
for i := 3; i < len(rowIDs); i += len(rowKeysBase) {
|
||||
delete(expected, colKeys[i])
|
||||
}
|
||||
}
|
||||
// and then we break the mutex and close the Qcxs
|
||||
for id, node := range nodesByID {
|
||||
field, err := node.API.Field(ctx, indexData.indexName, fieldData.fieldName)
|
||||
if err != nil {
|
||||
t.Fatalf("requesting field %s from node %s: %v", fieldData.fieldName, id, err)
|
||||
}
|
||||
pilosa.CorruptAMutex(t, field, qcxsByID[id])
|
||||
err = qcxsByID[id].Finish()
|
||||
if err != nil {
|
||||
t.Fatalf("closing out transaction on node %s: %v", id, err)
|
||||
}
|
||||
}
|
||||
qcx := m0.API.Txf().NewQcx()
|
||||
defer qcx.Abort()
|
||||
|
||||
results, err := m0.API.MutexCheck(ctx, qcx, indexData.indexName, fieldData.fieldName)
|
||||
if err != nil {
|
||||
t.Fatalf("checking mutexes: %v", err)
|
||||
}
|
||||
if keyedField {
|
||||
// this just sorta comes out this way with our hashing; these are
|
||||
// the things which were in position 1 of their shards, and did
|
||||
// not have a value which happens to map to 3.
|
||||
mapped, ok := results.(map[string][]string)
|
||||
if !ok {
|
||||
t.Fatalf("expected map[string][]string, got %T", results)
|
||||
}
|
||||
seen := 0
|
||||
for k, v := range mapped {
|
||||
seen++
|
||||
if _, ok := expected[k]; !ok {
|
||||
t.Fatalf("unexpected collision on key %q", k)
|
||||
}
|
||||
if len(v) != 2 {
|
||||
t.Fatalf("expected exactly two collisions")
|
||||
}
|
||||
}
|
||||
if seen != len(expected) {
|
||||
t.Fatalf("expected exactly %d records to have collisions, got %d", len(expected), seen)
|
||||
}
|
||||
} else {
|
||||
mapped, ok := results.(map[string][]uint64)
|
||||
if !ok {
|
||||
t.Fatalf("expected map[string][]uint64, got %T", results)
|
||||
}
|
||||
seen := 0
|
||||
for k, v := range mapped {
|
||||
seen++
|
||||
if _, ok := expected[k]; !ok {
|
||||
t.Fatalf("unexpected collision on key %q", k)
|
||||
}
|
||||
if len(v) != 2 {
|
||||
t.Fatalf("expected exactly two collisions")
|
||||
}
|
||||
}
|
||||
if seen != len(expected) {
|
||||
t.Fatalf("expected exactly %d records to have collisions, got %d", len(expected), seen)
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -84,6 +84,8 @@ type InternalClient interface {
|
|||
|
||||
GetNodeUsage(ctx context.Context, uri *URI) (map[string]NodeUsage, error)
|
||||
GetPastQueries(ctx context.Context, uri *URI) ([]PastQueryStatus, error)
|
||||
|
||||
MutexCheck(ctx context.Context, uri *URI, index string, field string) (map[uint64]map[uint64][]uint64, error)
|
||||
}
|
||||
|
||||
//===============
|
||||
|
|
@ -251,3 +253,7 @@ func (n nopInternalClient) GetNodeUsage(ctx context.Context, uri *URI) (map[stri
|
|||
func (n nopInternalClient) GetPastQueries(ctx context.Context, uri *URI) ([]PastQueryStatus, error) {
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
func (n nopInternalClient) MutexCheck(ctx context.Context, uri *URI, index, field string) (map[uint64]map[uint64][]uint64, error) {
|
||||
return nil, nil
|
||||
}
|
||||
|
|
|
|||
17
field.go
17
field.go
|
|
@ -1229,6 +1229,23 @@ func (f *Field) Row(tx Tx, rowID uint64) (*Row, error) {
|
|||
}
|
||||
}
|
||||
|
||||
// mutexCheck performs a sanity-check on the available fragments for a
|
||||
// field. The return is map[column]map[shard][]values for collisions only.
|
||||
func (f *Field) MutexCheck(ctx context.Context, qcx *Qcx) (map[uint64]map[uint64][]uint64, error) {
|
||||
if f.Type() != FieldTypeMutex {
|
||||
return nil, errors.New("mutex check only valid for mutex fields")
|
||||
}
|
||||
f.mu.RLock()
|
||||
defer f.mu.RUnlock()
|
||||
standard := f.viewMap[viewStandard]
|
||||
if standard == nil {
|
||||
// no standard view present means we've never needed to create it,
|
||||
// so it has no bits set, so it has no extra bits set.
|
||||
return nil, nil
|
||||
}
|
||||
return standard.mutexCheck(ctx, qcx)
|
||||
}
|
||||
|
||||
// SetBit sets a bit on a view within the field.
|
||||
func (f *Field) SetBit(tx Tx, rowID, colID uint64, t *time.Time) (changed bool, err error) {
|
||||
viewName := viewStandard
|
||||
|
|
|
|||
|
|
@ -27,6 +27,7 @@ import (
|
|||
|
||||
"github.com/pilosa/pilosa/v2/pql"
|
||||
"github.com/pilosa/pilosa/v2/roaring"
|
||||
"github.com/pilosa/pilosa/v2/shardwidth"
|
||||
"github.com/pilosa/pilosa/v2/testhook"
|
||||
)
|
||||
|
||||
|
|
@ -922,3 +923,27 @@ func TestBSIGroup_TxReopenDB(t *testing.T) {
|
|||
// the test: can we re-open a BSI fragment under Tx store
|
||||
_ = f.Reopen()
|
||||
}
|
||||
|
||||
func CorruptAMutex(tb testing.TB, field *Field, qcx *Qcx) {
|
||||
v := field.view(viewStandard)
|
||||
if v == nil {
|
||||
tb.Fatalf("creating view failed")
|
||||
}
|
||||
frags := v.allFragments()
|
||||
for _, frag := range frags {
|
||||
func() {
|
||||
tx, finisher, err := qcx.GetTx(Txo{Write: true, Index: field.idx, Shard: frag.shard})
|
||||
defer finisher(&err)
|
||||
if err != nil {
|
||||
tb.Fatalf("getting tx: %v", err)
|
||||
}
|
||||
// set a bonus bit, bypassing the mutex handling
|
||||
frag.mu.Lock()
|
||||
_, err = frag.unprotectedSetBit(tx, 3, (frag.shard<<shardwidth.Exponent)+1)
|
||||
frag.mu.Unlock()
|
||||
if err != nil {
|
||||
tb.Fatalf("setting bit: %v", err)
|
||||
}
|
||||
}()
|
||||
}
|
||||
}
|
||||
|
|
|
|||
11
fragment.go
11
fragment.go
|
|
@ -585,6 +585,17 @@ func (f *fragment) closeStorage() error {
|
|||
return nil
|
||||
}
|
||||
|
||||
// mutexCheck checks for any entries in fragment which violate the mutex
|
||||
// property of having only one value set for a given column ID.
|
||||
func (f *fragment) mutexCheck(tx Tx) (map[uint64][]uint64, error) {
|
||||
dup := roaring.NewBitmapMutexDupFilter(f.shard << shardwidth.Exponent)
|
||||
err := tx.ApplyFilter(f.index(), f.field(), f.view(), f.shard, 0, dup)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return dup.Report(), nil
|
||||
}
|
||||
|
||||
// row returns a row by ID.
|
||||
func (f *fragment) row(tx Tx, rowID uint64) (*Row, error) {
|
||||
f.mu.Lock()
|
||||
|
|
|
|||
14
go.sum
14
go.sum
|
|
@ -65,7 +65,6 @@ github.com/envoyproxy/go-control-plane v0.9.0/go.mod h1:YTl/9mNaCwkRvm6d1a2C3ymF
|
|||
github.com/envoyproxy/go-control-plane v0.9.4/go.mod h1:6rpuAdCZL397s3pYoYcLgu1mIlRU8Am5FuJP05cCM98=
|
||||
github.com/envoyproxy/protoc-gen-validate v0.1.0/go.mod h1:iSmxcyjqTsJpI2R4NaDN7+kN2VEUnK/pcBlmesArF7c=
|
||||
github.com/fatih/color v1.7.0/go.mod h1:Zm6kSWBoL9eyXnKyktHP6abPY2pDugNf5KwzbycvMj4=
|
||||
github.com/fsnotify/fsnotify v1.4.7 h1:IXs+QLmnXW2CcXuY+8Mzv/fWEsPGWxqefPtCP5CnV9I=
|
||||
github.com/fsnotify/fsnotify v1.4.7/go.mod h1:jwhsz4b93w/PPRr/qN1Yymfu8t87LnFCMoQvtojpjFo=
|
||||
github.com/fsnotify/fsnotify v1.4.9 h1:hsms1Qyu0jgnwNXIxa+/V/PDsU6CfLf6CNO8H7IWoS4=
|
||||
github.com/fsnotify/fsnotify v1.4.9/go.mod h1:znqG4EE+3YCdAaPaxE2ZRY/06pZUdp0tY4IgpuI1SZQ=
|
||||
|
|
@ -95,7 +94,6 @@ github.com/golang/mock v1.3.1/go.mod h1:sBzyDLLjw3U8JLTeZvSv8jJB+tU5PVekmnlKIyFU
|
|||
github.com/golang/protobuf v1.2.0/go.mod h1:6lQm79b+lXiMfvg/cZm0SGofjICqVBUtrP5yJMmIC1U=
|
||||
github.com/golang/protobuf v1.3.1/go.mod h1:6lQm79b+lXiMfvg/cZm0SGofjICqVBUtrP5yJMmIC1U=
|
||||
github.com/golang/protobuf v1.3.2/go.mod h1:6lQm79b+lXiMfvg/cZm0SGofjICqVBUtrP5yJMmIC1U=
|
||||
github.com/golang/protobuf v1.3.3 h1:gyjaxf+svBWX08ZjK86iN9geUJF0H6gp2IRKX6Nf6/I=
|
||||
github.com/golang/protobuf v1.3.3/go.mod h1:vzj43D7+SQXF/4pzW/hwtAqwc6iTitCiVSaWz5lYuqw=
|
||||
github.com/golang/protobuf v1.4.0-rc.1/go.mod h1:ceaxUfeHdC40wWswd/P6IGgMaK3YpKi5j83Wpe3EHw8=
|
||||
github.com/golang/protobuf v1.4.0-rc.1.0.20200221234624-67d41d38c208/go.mod h1:xKAWHe0F5eneWXFV3EuXVDTCmh+JuBKY0li0aMyXATA=
|
||||
|
|
@ -110,7 +108,6 @@ github.com/google/btree v1.0.0/go.mod h1:lNA+9X1NB3Zf8V7Ke586lFgjr2dZNuvo3lPJSGZ
|
|||
github.com/google/go-cmp v0.2.0/go.mod h1:oXzfMopK8JAjlY9xF4vHSVASa0yLyX7SntLO5aqRK0M=
|
||||
github.com/google/go-cmp v0.3.0/go.mod h1:8QqcDgzrUqlUb/G2PQTWiueGozuR1884gddMywk6iLU=
|
||||
github.com/google/go-cmp v0.3.1/go.mod h1:8QqcDgzrUqlUb/G2PQTWiueGozuR1884gddMywk6iLU=
|
||||
github.com/google/go-cmp v0.4.0 h1:xsAVV57WRhGj6kEIi8ReJzQlHHqcBYCElAvkovg3B/4=
|
||||
github.com/google/go-cmp v0.4.0/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE=
|
||||
github.com/google/go-cmp v0.5.2 h1:X2ev0eStA3AbceY54o37/0PQ/UWqKEiiO2dKL5OPaFM=
|
||||
github.com/google/go-cmp v0.5.2/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE=
|
||||
|
|
@ -147,12 +144,10 @@ github.com/hashicorp/go-rootcerts v1.0.0/go.mod h1:K6zTfqpRlCUIjkwsN4Z+hiSfzSTQa
|
|||
github.com/hashicorp/go-sockaddr v1.0.0 h1:GeH6tui99pF4NJgfnhp+L6+FfobzVW3Ah46sLo0ICXs=
|
||||
github.com/hashicorp/go-sockaddr v1.0.0/go.mod h1:7Xibr9yA9JjQq1JpNB2Vw7kxv8xerXegt+ozgdvDeDU=
|
||||
github.com/hashicorp/go-syslog v1.0.0/go.mod h1:qPfqrKkXGihmCqbJM2mZgkZGvKG1dFdvsLplgctolz4=
|
||||
github.com/hashicorp/go-uuid v1.0.0 h1:RS8zrF7PhGwyNPOtxSClXXj9HA8feRnJzgnI1RJCSnM=
|
||||
github.com/hashicorp/go-uuid v1.0.0/go.mod h1:6SBZvOh/SIDV7/2o3Jml5SYk/TvGqwFJ/bN7x4byOro=
|
||||
github.com/hashicorp/go-uuid v1.0.1 h1:fv1ep09latC32wFoVwnqcnKJGnMSdBanPczbHAYm1BE=
|
||||
github.com/hashicorp/go-uuid v1.0.1/go.mod h1:6SBZvOh/SIDV7/2o3Jml5SYk/TvGqwFJ/bN7x4byOro=
|
||||
github.com/hashicorp/go.net v0.0.1/go.mod h1:hjKkEWcCURg++eb33jQU7oqQcI9XDCnUzHA0oac0k90=
|
||||
github.com/hashicorp/golang-lru v0.5.0 h1:CL2msUPvZTLb5O648aiLNJw3hnBxN2+1Jq8rCOH9wdo=
|
||||
github.com/hashicorp/golang-lru v0.5.0/go.mod h1:/m3WP610KZHVQ1SGc6re/UDhFvYD7pJ4Ao+sR/qLZy8=
|
||||
github.com/hashicorp/golang-lru v0.5.1 h1:0hERBMJE1eitiLkihrMvRVBYAkpHzc/J3QdDN+dAcgU=
|
||||
github.com/hashicorp/golang-lru v0.5.1/go.mod h1:/m3WP610KZHVQ1SGc6re/UDhFvYD7pJ4Ao+sR/qLZy8=
|
||||
|
|
@ -172,15 +167,12 @@ github.com/jtolds/gls v4.20.0+incompatible h1:xdiiI2gbIgH/gLH7ADydsJ1uDOEzR8yvV7
|
|||
github.com/jtolds/gls v4.20.0+incompatible/go.mod h1:QJZ7F/aHp+rZTRtaJ1ow/lLfFfVYBRgL+9YlvaHOwJU=
|
||||
github.com/julienschmidt/httprouter v1.2.0/go.mod h1:SYymIcj16QtmaHHD7aYtjjsJG7VTCxuUUipMqKk8s4w=
|
||||
github.com/kisielk/errcheck v1.1.0/go.mod h1:EZBBE59ingxPouuu3KfxchcWSUPOHkagtvWXihfKN4Q=
|
||||
github.com/kisielk/gotool v1.0.0 h1:AV2c/EiW3KqPNT9ZKl07ehoAGi4C5/01Cfbblndcapg=
|
||||
github.com/kisielk/gotool v1.0.0/go.mod h1:XhKaO+MFFWcvkIS/tQcRk01m1F5IRFswLeQ+oQHNcck=
|
||||
github.com/konsorten/go-windows-terminal-sequences v1.0.1/go.mod h1:T0+1ngSBFLxvqU3pZ+m/2kptfBszLMUkC4ZK/EgS/cQ=
|
||||
github.com/konsorten/go-windows-terminal-sequences v1.0.2/go.mod h1:T0+1ngSBFLxvqU3pZ+m/2kptfBszLMUkC4ZK/EgS/cQ=
|
||||
github.com/kr/logfmt v0.0.0-20140226030751-b84e30acd515/go.mod h1:+0opPa2QZZtGFBFZlji/RkVcI2GknAs/DXo4wKdlNEc=
|
||||
github.com/kr/pretty v0.1.0 h1:L/CwN0zerZDmRFUapSPitk6f+Q3+0za1rQkzVuMiMFI=
|
||||
github.com/kr/pretty v0.1.0/go.mod h1:dAy3ld7l9f0ibDNOQOHHMYYIIbhfbHSm3C4ZsoJORNo=
|
||||
github.com/kr/pty v1.1.1/go.mod h1:pFQYn66WHrOpPYNljwOMqo10TkYh1fy3cYio2l3bCsQ=
|
||||
github.com/kr/text v0.1.0 h1:45sCR5RtlFHMR4UwH9sdQ5TC8v0qDQCHnXt+kaKSTVE=
|
||||
github.com/kr/text v0.1.0/go.mod h1:4Jbv+DJW3UT/LiOwJeYQe1efqtUx/iVham/4vfdArNI=
|
||||
github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY=
|
||||
github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE=
|
||||
|
|
@ -196,7 +188,6 @@ github.com/miekg/dns v1.0.14 h1:9jZdLNd/P4+SfEJ0TNyxYpsK8N4GtfylBLqtbYN1sbA=
|
|||
github.com/miekg/dns v1.0.14/go.mod h1:W1PPwlIAgtquWBMBEV9nkV9Cazfe8ScdGz/Lj7v3Nrg=
|
||||
github.com/mitchellh/cli v1.0.0/go.mod h1:hNIlj7HEI86fIcpObd7a0FcrxTWetlwJDGcceTlRvqc=
|
||||
github.com/mitchellh/go-homedir v1.0.0/go.mod h1:SfyaCUpYCn1Vlf4IUYiD9fPX4A5wJrkLzIz1N1q0pr0=
|
||||
github.com/mitchellh/go-homedir v1.1.0 h1:lukF9ziXFxDFPkA1vsr5zpc1XuPDn/wFntq5mG+4E0Y=
|
||||
github.com/mitchellh/go-homedir v1.1.0/go.mod h1:SfyaCUpYCn1Vlf4IUYiD9fPX4A5wJrkLzIz1N1q0pr0=
|
||||
github.com/mitchellh/go-testing-interface v1.0.0/go.mod h1:kRemZodwjscx+RGhAo8eIhFbs2+BFgRtFPeD/KE+zxI=
|
||||
github.com/mitchellh/gox v0.4.0/go.mod h1:Sd9lOJ0+aimLBi73mGofS1ycjY8lL3uZM3JPS42BGNg=
|
||||
|
|
@ -291,11 +282,9 @@ github.com/spf13/viper v1.7.0/go.mod h1:8WkrPz2fc9jxqZNCJI/76HCieCp4Q8HaLFoCha5q
|
|||
github.com/spf13/viper v1.7.1 h1:pM5oEahlgWv/WnHXpgbKz7iLIxRf65tye2Ci+XFK5sk=
|
||||
github.com/spf13/viper v1.7.1/go.mod h1:8WkrPz2fc9jxqZNCJI/76HCieCp4Q8HaLFoCha5qpdg=
|
||||
github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
|
||||
github.com/stretchr/objx v0.1.1 h1:2vfRuCMp5sSVIDSqO8oNnWJq7mPa6KVP3iPIwFBuy8A=
|
||||
github.com/stretchr/objx v0.1.1/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
|
||||
github.com/stretchr/testify v1.2.2/go.mod h1:a8OnRcib4nhh0OaRAV+Yts87kKdq0PP7pXfy6kDkUVs=
|
||||
github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI=
|
||||
github.com/stretchr/testify v1.4.0 h1:2E4SXV/wtOkTonXsotYi4li6zVWxYlZuYNCXe9XRJyk=
|
||||
github.com/stretchr/testify v1.4.0/go.mod h1:j7eGeouHqKxXV5pUuKE4zz7dFj8WfuZ+81PSLYec5m4=
|
||||
github.com/stretchr/testify v1.6.1 h1:hDPOHmpOpP40lSULcqw7IrRb/u7w6RpDC9399XyoNd0=
|
||||
github.com/stretchr/testify v1.6.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg=
|
||||
|
|
@ -401,7 +390,6 @@ golang.org/x/sys v0.0.0-20191001151750-bb3f8db39f24/go.mod h1:h1NjWce9XRLGQEsW7w
|
|||
golang.org/x/sys v0.0.0-20191005200804-aed5e4c7ecf9/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
|
||||
golang.org/x/sys v0.0.0-20191220142924-d4481acd189f/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
|
||||
golang.org/x/sys v0.0.0-20200202164722-d101bd2416d5/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
|
||||
golang.org/x/sys v0.0.0-20200323222414-85ca7c5b95cd h1:xhmwyvizuTgC2qz7ZlMluP20uW+C3Rm0FD/WLDX8884=
|
||||
golang.org/x/sys v0.0.0-20200323222414-85ca7c5b95cd/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
|
||||
golang.org/x/sys v0.0.0-20201024232916-9f70ab9862d5/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
|
||||
golang.org/x/sys v0.0.0-20201214095126-aec9a390925b h1:tv7/y4pd+sR8bcNb2D6o7BNU6zjWm0VjQLac+w7fNNM=
|
||||
|
|
@ -453,7 +441,6 @@ google.golang.org/genproto v0.0.0-20190418145605-e7d98fc518a7/go.mod h1:VzzqZJRn
|
|||
google.golang.org/genproto v0.0.0-20190425155659-357c62f0e4bb/go.mod h1:VzzqZJRnGkLBvHegQrXjBqPurQTc5/KpmUdxsrq26oE=
|
||||
google.golang.org/genproto v0.0.0-20190502173448-54afdca5d873/go.mod h1:VzzqZJRnGkLBvHegQrXjBqPurQTc5/KpmUdxsrq26oE=
|
||||
google.golang.org/genproto v0.0.0-20190801165951-fa694d86fc64/go.mod h1:DMBHOl98Agz4BDEuKkezgsaosCRResVns1a3J2ZsMNc=
|
||||
google.golang.org/genproto v0.0.0-20190819201941-24fa4b261c55 h1:gSJIx1SDwno+2ElGhA4+qG2zF97qiUzTM+rQ0klBOcE=
|
||||
google.golang.org/genproto v0.0.0-20190819201941-24fa4b261c55/go.mod h1:DMBHOl98Agz4BDEuKkezgsaosCRResVns1a3J2ZsMNc=
|
||||
google.golang.org/genproto v0.0.0-20190911173649-1774047e7e51/go.mod h1:IbNlFCBrqXvoKpeg0TB2l7cyZUmoaFKYIwrEpbDKLA8=
|
||||
google.golang.org/genproto v0.0.0-20191108220845-16a3f7862a1a h1:Ob5/580gVHBJZgXnff1cZDbG+xLtMVE5mDRTe+nIsX4=
|
||||
|
|
@ -483,7 +470,6 @@ gopkg.in/ini.v1 v1.51.0/go.mod h1:pNLf8WUiyNEtQjuu5G5vTm06TEv9tsIgeAvK8hOrP4k=
|
|||
gopkg.in/resty.v1 v1.12.0/go.mod h1:mDo4pnntr5jdWRML875a/NmxYqAlA73dVijT2AXvQQo=
|
||||
gopkg.in/yaml.v2 v2.0.0-20170812160011-eb3733d160e7/go.mod h1:JAlM8MvJe8wmxCU4Bli9HhUf9+ttbYbLASfIpnQbh74=
|
||||
gopkg.in/yaml.v2 v2.2.1/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI=
|
||||
gopkg.in/yaml.v2 v2.2.2 h1:ZCJp+EgiOT7lHqUV2J862kp8Qj64Jo6az82+3Td9dZw=
|
||||
gopkg.in/yaml.v2 v2.2.2/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI=
|
||||
gopkg.in/yaml.v2 v2.2.4/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI=
|
||||
gopkg.in/yaml.v2 v2.2.8/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI=
|
||||
|
|
|
|||
|
|
@ -133,6 +133,34 @@ func (c *InternalClient) Schema(ctx context.Context) ([]*pilosa.IndexInfo, error
|
|||
return rsp.Indexes, nil
|
||||
}
|
||||
|
||||
// MutexCheck uses the mutex-check endpoint to request mutex collision data
|
||||
// from a single node.
|
||||
func (c *InternalClient) MutexCheck(ctx context.Context, uri *pilosa.URI, indexName string, fieldName string) (map[uint64]map[uint64][]uint64, error) {
|
||||
if uri == nil {
|
||||
uri = c.defaultURI
|
||||
}
|
||||
u := uri.Path(fmt.Sprintf("/internal/index/%s/field/%s/mutex-check", indexName, fieldName))
|
||||
req, err := http.NewRequest("GET", u, nil)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "creating request")
|
||||
}
|
||||
req.Header.Set("Accept", "application/json")
|
||||
req.Header.Set("User-Agent", "pilosa/"+pilosa.Version)
|
||||
|
||||
resp, err := c.executeRequest(req.WithContext(ctx))
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "executing request")
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
return nil, errors.Errorf("unexpected status code: %s", resp.Status)
|
||||
}
|
||||
var out map[uint64]map[uint64][]uint64
|
||||
dec := json.NewDecoder(resp.Body)
|
||||
err = dec.Decode(&out)
|
||||
return out, err
|
||||
}
|
||||
|
||||
func (c *InternalClient) PostSchema(ctx context.Context, uri *pilosa.URI, s *pilosa.Schema, remote bool) error {
|
||||
u := uri.Path(fmt.Sprintf("/schema?remote=%v", remote))
|
||||
buf, err := json.Marshal(s)
|
||||
|
|
|
|||
|
|
@ -386,6 +386,7 @@ func newRouter(handler *Handler) http.Handler {
|
|||
router.HandleFunc("/index/{index}/field/{field}", handler.handlePostField).Methods("POST").Name("PostField")
|
||||
router.HandleFunc("/index/{index}/field/{field}", handler.handleDeleteField).Methods("DELETE").Name("DeleteField")
|
||||
router.HandleFunc("/index/{index}/field/{field}/import", handler.handlePostImport).Methods("POST").Name("PostImport")
|
||||
router.HandleFunc("/index/{index}/field/{field}/mutex-check", handler.handleGetMutexCheck).Methods("GET").Name("GetMutexCheck")
|
||||
router.HandleFunc("/index/{index}/field/{field}/import-roaring/{shard}", handler.handlePostImportRoaring).Methods("POST").Name("PostImportRoaring")
|
||||
router.HandleFunc("/index/{index}/query", handler.handlePostQuery).Methods("POST").Name("PostQuery")
|
||||
router.HandleFunc("/info", handler.handleGetInfo).Methods("GET").Name("GetInfo")
|
||||
|
|
@ -422,6 +423,7 @@ func newRouter(handler *Handler) http.Handler {
|
|||
router.HandleFunc("/internal/translate/data", handler.handlePostTranslateData).Methods("POST").Name("PostTranslateData")
|
||||
router.HandleFunc("/internal/translate/keys", handler.handlePostTranslateKeys).Methods("POST").Name("PostTranslateKeys")
|
||||
router.HandleFunc("/internal/translate/ids", handler.handlePostTranslateIDs).Methods("POST").Name("PostTranslateIDs")
|
||||
router.HandleFunc("/internal/index/{index}/field/{field}/mutex-check", handler.handleInternalGetMutexCheck).Methods("GET").Name("InternalGetMutexCheck")
|
||||
router.HandleFunc("/internal/index/{index}/field/{field}/attr/diff", handler.handlePostFieldAttrDiff).Methods("POST").Name("PostFieldAttrDiff")
|
||||
router.HandleFunc("/internal/index/{index}/field/{field}/remote-available-shards/{shardID}", handler.handleDeleteRemoteAvailableShard).Methods("DELETE")
|
||||
router.HandleFunc("/internal/nodes", handler.handleGetNodes).Methods("GET").Name("GetNodes")
|
||||
|
|
@ -2479,6 +2481,58 @@ func (h *Handler) handlePostImportColumnAttrs(w http.ResponseWriter, r *http.Req
|
|||
}
|
||||
}
|
||||
|
||||
// handleGetMutexCheck handles /mutex-check requests.
|
||||
func (h *Handler) handleGetMutexCheck(w http.ResponseWriter, r *http.Request) {
|
||||
if !validHeaderAcceptJSON(r.Header) {
|
||||
http.Error(w, "JSON only acceptable response", http.StatusNotAcceptable)
|
||||
return
|
||||
}
|
||||
// Get index and field type to determine how to handle the
|
||||
// import data.
|
||||
indexName, fieldName := mux.Vars(r)["index"], mux.Vars(r)["field"]
|
||||
qcx := h.api.Txf().NewQcx()
|
||||
defer qcx.Abort()
|
||||
out, err := h.api.MutexCheck(r.Context(), qcx, indexName, fieldName)
|
||||
if err != nil {
|
||||
http.Error(w, err.Error(), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
outBytes, err := json.Marshal(out)
|
||||
if err != nil {
|
||||
http.Error(w, fmt.Sprintf("marshalling response: %v", err), http.StatusInternalServerError)
|
||||
}
|
||||
_, err = w.Write(outBytes)
|
||||
if err != nil {
|
||||
h.logger.Printf("writing mutex-check response: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// handleInternalGetMutexCheck handles internal (non-forwarding )/mutex-check requests.
|
||||
func (h *Handler) handleInternalGetMutexCheck(w http.ResponseWriter, r *http.Request) {
|
||||
if !validHeaderAcceptJSON(r.Header) {
|
||||
http.Error(w, "JSON only acceptable response", http.StatusNotAcceptable)
|
||||
return
|
||||
}
|
||||
// Get index and field type to determine how to handle the
|
||||
// import data.
|
||||
indexName, fieldName := mux.Vars(r)["index"], mux.Vars(r)["field"]
|
||||
qcx := h.api.Txf().NewQcx()
|
||||
defer qcx.Abort()
|
||||
out, err := h.api.MutexCheckNode(r.Context(), qcx, indexName, fieldName)
|
||||
if err != nil {
|
||||
http.Error(w, err.Error(), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
outBytes, err := json.Marshal(out)
|
||||
if err != nil {
|
||||
http.Error(w, fmt.Sprintf("marshalling response: %v", err), http.StatusInternalServerError)
|
||||
}
|
||||
_, err = w.Write(outBytes)
|
||||
if err != nil {
|
||||
h.logger.Printf("writing mutex-check response: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// handlePostImportRoaring
|
||||
func (h *Handler) handlePostImportRoaring(w http.ResponseWriter, r *http.Request) {
|
||||
// Verify that request is only communicating over protobufs.
|
||||
|
|
|
|||
|
|
@ -755,6 +755,75 @@ func NewBitmapRangeFilter(min, max FilterKey, keyCallback func(FilterKey, int32)
|
|||
return &BitmapRangeFilter{min: min, max: max, kcb: keyCallback, dcb: dataCallback}
|
||||
}
|
||||
|
||||
// BitmapMutexDupFilter is a filter which identifies cases where the same
|
||||
// position has a bit set in more than one row.
|
||||
//
|
||||
// We keep a slice of the first value seen for every row, with ^0 as the
|
||||
// default; when that's already set, things get appended to the entries in
|
||||
// the map. At the end, for each entry in the map, we also add its first
|
||||
// value to it. Thus, the map holds all the entries, but we're only using
|
||||
// the map in the (hopefully rarer) cases where there's duplicate values.
|
||||
//
|
||||
// The slice is local-coordinates (first column 0), but the map is global
|
||||
// coordinates (first column is whatever base was).
|
||||
type BitmapMutexDupFilter struct {
|
||||
base uint64 // the offset of 0 for this, used to accommodate shard offsets
|
||||
extra map[uint64][]uint64 // extra values observed
|
||||
first []uint64 // first values observed
|
||||
}
|
||||
|
||||
var _ BitmapFilter = &BitmapMutexDupFilter{}
|
||||
|
||||
func NewBitmapMutexDupFilter(base uint64) *BitmapMutexDupFilter {
|
||||
filter := &BitmapMutexDupFilter{
|
||||
base: base,
|
||||
extra: map[uint64][]uint64{},
|
||||
first: make([]uint64, 1<<shardwidth.Exponent),
|
||||
}
|
||||
for i := range filter.first {
|
||||
filter.first[i] = ^uint64(0)
|
||||
}
|
||||
return filter
|
||||
}
|
||||
|
||||
func (b *BitmapMutexDupFilter) ConsiderKey(key FilterKey, n int32) FilterResult {
|
||||
if n > 0 {
|
||||
return key.NeedData()
|
||||
}
|
||||
return key.RejectOne()
|
||||
}
|
||||
|
||||
func (b *BitmapMutexDupFilter) ConsiderData(key FilterKey, data *Container) FilterResult {
|
||||
value, basePos := uint64(key)>>rowExponent, uint64(key&keyMask)<<16
|
||||
containerCallback(data, func(u uint16) {
|
||||
pos := basePos + uint64(u)
|
||||
if b.first[pos] != ^uint64(0) {
|
||||
b.extra[pos+b.base] = append(b.extra[pos+b.base], value)
|
||||
} else {
|
||||
b.first[pos] = value
|
||||
}
|
||||
})
|
||||
return key.MatchOne()
|
||||
}
|
||||
|
||||
// Report returns the set of duplicate values identified.
|
||||
func (b *BitmapMutexDupFilter) Report() map[uint64][]uint64 {
|
||||
// copy values into extra, and remove them from first, so calling
|
||||
// Report() again won't cause double-appends.
|
||||
for k, v := range b.extra {
|
||||
kpos := k % (1 << shardwidth.Exponent)
|
||||
if b.first[kpos] != ^uint64(0) {
|
||||
v = append(v, 0)
|
||||
// prepend so the lowest value goes at the beginning
|
||||
copy(v[1:], v[:])
|
||||
v[0] = b.first[kpos]
|
||||
b.first[kpos] = ^uint64(0)
|
||||
b.extra[k] = v
|
||||
}
|
||||
}
|
||||
return b.extra
|
||||
}
|
||||
|
||||
// ApplyFilterToIterator is a simplistic implementation that applies a bitmap
|
||||
// filter to a ContainerIterator, returning an error if it encounters an error.
|
||||
//
|
||||
|
|
|
|||
|
|
@ -349,3 +349,50 @@ func TestFilterWithRows(t *testing.T) {
|
|||
}
|
||||
|
||||
}
|
||||
|
||||
func TestMutexDupFilter(t *testing.T) {
|
||||
tests := []struct{
|
||||
pairs [][2]uint64
|
||||
expect map[uint64][]uint64
|
||||
}{
|
||||
{
|
||||
pairs: [][2]uint64{{0, 0}, {1, 0}, {0, 1}},
|
||||
expect: map[uint64][]uint64{0: {0, 1}},
|
||||
},
|
||||
{
|
||||
pairs: [][2]uint64{{0, 0}, {1, 0}, {0, 1}, {0, 2}},
|
||||
expect: map[uint64][]uint64{0: {0, 1, 2}},
|
||||
},
|
||||
}
|
||||
for num, test := range tests {
|
||||
t.Run(fmt.Sprintf("case%d", num), func(t *testing.T) {
|
||||
b := NewSliceBitmap()
|
||||
for _, p := range test.pairs {
|
||||
v := (p[1] << shardwidth.Exponent) | p[0]
|
||||
b.DirectAdd(v)
|
||||
}
|
||||
dup := NewBitmapMutexDupFilter(0)
|
||||
iter, _ := b.Containers.Iterator(0)
|
||||
err := ApplyFilterToIterator(dup, iter)
|
||||
if err != nil {
|
||||
t.Fatalf("applying filter: %v", err)
|
||||
}
|
||||
expected := test.expect
|
||||
got := dup.Report()
|
||||
if len(expected) != len(got) {
|
||||
t.Fatalf("expected %d entries in duplicate map, got %d", len(expected), len(got))
|
||||
}
|
||||
for k, v := range expected {
|
||||
gv := got[k]
|
||||
if len(v) != len(gv) {
|
||||
t.Fatalf("for id %d, expected %d (len %d), got %d (len %d)", k, v, len(v), gv, len(gv))
|
||||
}
|
||||
for j := range v {
|
||||
if gv[j] != v[j] {
|
||||
t.Fatalf("for id %d, expected %d, got %d", k, v[j], gv[j])
|
||||
}
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -4448,10 +4448,13 @@ func intersectionCallbackArrayArray(a, b *Container, fn func(uint16)) {
|
|||
if (na << 2) < nb {
|
||||
for _, va := range ca {
|
||||
for cb[0] < va {
|
||||
if len(cb) > 8 && cb[0] < va {
|
||||
// try to skip ahead a bit faster
|
||||
for len(cb) > 7 && cb[7] < va {
|
||||
cb = cb[8:]
|
||||
}
|
||||
cb = cb[1:]
|
||||
for len(cb) > 0 && cb[0] < va {
|
||||
cb = cb[1:]
|
||||
}
|
||||
if len(cb) == 0 {
|
||||
return
|
||||
}
|
||||
|
|
@ -4540,7 +4543,7 @@ func intersectionCallbackBitmapRun(a, b *Container, fn func(uint16)) {
|
|||
}
|
||||
}
|
||||
|
||||
func intersectionCallbackArrayBitmap(a, b *Container, fn func(uint16)) (n int32) {
|
||||
func intersectionCallbackArrayBitmap(a, b *Container, fn func(uint16)) {
|
||||
statsHit("intersectionCount/ArrayBitmap")
|
||||
bitmap := b.bitmap()
|
||||
ln := len(bitmap)
|
||||
|
|
@ -4550,9 +4553,10 @@ func intersectionCallbackArrayBitmap(a, b *Container, fn func(uint16)) (n int32)
|
|||
break
|
||||
}
|
||||
off := val % 64
|
||||
n += int32(bitmap[i]>>off) & 1
|
||||
if (bitmap[i]>>off)&1 != 0 {
|
||||
fn(val)
|
||||
}
|
||||
}
|
||||
return n
|
||||
}
|
||||
|
||||
func intersectionCallbackBitmapBitmap(a, b *Container, fn func(uint16)) {
|
||||
|
|
|
|||
44
view.go
44
view.go
|
|
@ -447,6 +447,50 @@ func (v *view) row(txOrig Tx, rowID uint64) (*Row, error) {
|
|||
|
||||
}
|
||||
|
||||
// mutexCheck checks all available fragments for duplicate values. The return
|
||||
// is map[column]map[shard][]values for collisions only.
|
||||
func (v *view) mutexCheck(ctx context.Context, qcx *Qcx) (map[uint64]map[uint64][]uint64, error) {
|
||||
// We don't need the context, we just want the context-awareness on the error groups.
|
||||
// It would be nice if the inner functions could use this too...
|
||||
eg, _ := errgroup.WithContext(ctx)
|
||||
throttle := make(chan struct{}, runtime.NumCPU())
|
||||
frags := v.allFragments()
|
||||
results := make([]map[uint64][]uint64, len(frags))
|
||||
for i, frag := range frags {
|
||||
// local copies for the goroutine to use
|
||||
i, frag := i, frag
|
||||
eg.Go(func() error {
|
||||
// limit simultaneous parallel goroutines associated with this
|
||||
throttle <- struct{}{}
|
||||
defer func() {
|
||||
<-throttle
|
||||
}()
|
||||
tx, finisher, err := qcx.GetTx(Txo{Index: v.idx, Shard: frag.shard})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer finisher(&err)
|
||||
results[i], err = frag.mutexCheck(tx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
})
|
||||
}
|
||||
err := eg.Wait()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out := map[uint64]map[uint64][]uint64{}
|
||||
for i, result := range results {
|
||||
if len(result) == 0 {
|
||||
continue
|
||||
}
|
||||
out[frags[i].shard] = result
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// setBit sets a bit within the view.
|
||||
func (v *view) setBit(txOrig Tx, rowID, columnID uint64) (changed bool, err error) {
|
||||
shard := columnID / ShardWidth
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue