WIP: implement Schemator

This commit is contained in:
Travis 2021-02-07 23:44:26 -06:00
parent 114f6a8751
commit 3542134100
No known key found for this signature in database
GPG key ID: 37080CC2042BA34E
7 changed files with 371 additions and 126 deletions

88
api.go
View file

@ -57,6 +57,7 @@ type API struct {
importWork chan importJob
Serializer Serializer
schemator disco.Schemator
}
func (api *API) Holder() *Holder {
@ -72,6 +73,7 @@ func OptAPIServer(s *Server) apiOption {
a.holder = s.holder
a.cluster = s.cluster
a.Serializer = s.serializer
a.schemator = s.schemator
return nil
}
}
@ -211,37 +213,19 @@ func (api *API) CreateIndex(ctx context.Context, indexName string, options Index
return nil, errors.Wrap(err, "validating api method")
}
// Create a snapshot of the cluster to use for node/partition calculations.
snap := topology.NewClusterSnapshot(api.cluster.noder, api.cluster.Hasher, api.cluster.ReplicaN)
if !snap.IsPrimaryFieldTranslationNode(api.Node().ID) {
if err := api.server.defaultClient.CreateIndex(ctx, indexName, options); err != nil {
return nil, errors.Wrap(err, "forwarding CreateIndex to coordinator")
}
return api.holder.Index(indexName), nil
// Populate the create index message.
cim := &CreateIndexMessage{
Index: indexName,
CreatedAt: timestamp(),
Meta: &options,
}
// Create index.
index, err := api.holder.CreateIndex(indexName, options)
index, err := api.holder.CreateIndexAndBroadcast(cim)
if err != nil {
return nil, errors.Wrap(err, "creating index")
}
createdAt := timestamp()
index.mu.Lock()
index.createdAt = createdAt
index.mu.Unlock()
// Send the create index message to all nodes.
err = api.server.SendSync(
&CreateIndexMessage{
Index: indexName,
CreatedAt: createdAt,
Meta: &options,
})
if err != nil {
return nil, errors.Wrap(err, "sending CreateIndex message")
}
api.holder.Stats.Count(MetricCreateIndex, 1, 1.0)
return index, nil
}
@ -272,6 +256,11 @@ func (api *API) DeleteIndex(ctx context.Context, indexName string) error {
return errors.Wrap(err, "validating api method")
}
// Delete the index from etcd as the system of record.
if err := api.schemator.DeleteIndex(ctx, indexName); err != nil {
return errors.Wrapf(err, "deleting index from etcd: %s", indexName)
}
// Delete index from the holder.
err := api.holder.DeleteIndex(indexName)
if err != nil {
@ -301,23 +290,10 @@ func (api *API) CreateField(ctx context.Context, indexName string, fieldName str
return nil, errors.Wrap(err, "validating api method")
}
// Apply functional options.
fo := FieldOptions{}
for _, opt := range opts {
err := opt(&fo)
if err != nil {
return nil, NewBadRequestError(errors.Wrap(err, "applying option"))
}
}
// Create a snapshot of the cluster to use for node/partition calculations.
snap := topology.NewClusterSnapshot(api.cluster.noder, api.cluster.Hasher, api.cluster.ReplicaN)
if !snap.IsPrimaryFieldTranslationNode(api.Node().ID) {
if err := api.server.defaultClient.CreateFieldWithOptions(ctx, indexName, fieldName, fo); err != nil {
return nil, errors.Wrap(err, "forwarding CreateField to coordinator")
}
return api.holder.Field(indexName, fieldName), nil
// Apply and validate functional options.
fo, err := newFieldOptions(opts...)
if err != nil {
return nil, NewBadRequestError(errors.Wrap(err, "applying option"))
}
// Find index.
@ -326,27 +302,20 @@ func (api *API) CreateField(ctx context.Context, indexName string, fieldName str
return nil, newNotFoundError(ErrIndexNotFound, indexName)
}
// Populate the create field message.
cfm := &CreateFieldMessage{
Index: indexName,
Field: fieldName,
CreatedAt: timestamp(),
Meta: fo,
}
// Create field.
field, err := index.CreateField(fieldName, opts...)
field, err := index.CreateFieldAndBroadcast(cfm)
if err != nil {
return nil, errors.Wrap(err, "creating field")
}
createdAt := timestamp()
field.mu.Lock()
field.createdAt = createdAt
field.mu.Unlock()
// Send the create field message to all nodes.
err = api.server.SendSync(&CreateFieldMessage{
Index: indexName,
Field: fieldName,
CreatedAt: createdAt,
Meta: &fo,
})
if err != nil {
api.server.logger.Printf("problem sending CreateField message: %s", err)
return nil, errors.Wrap(err, "sending CreateField message")
}
api.holder.Stats.CountWithCustomTags(MetricCreateField, 1, 1.0, []string{fmt.Sprintf("index:%s", indexName)})
return field, nil
}
@ -571,6 +540,11 @@ func (api *API) DeleteField(ctx context.Context, indexName string, fieldName str
return newNotFoundError(ErrIndexNotFound, indexName)
}
// Delete the field from etcd as the system of record.
if err := api.schemator.DeleteField(ctx, indexName, fieldName); err != nil {
return errors.Wrapf(err, "deleting field from etcd: %s/%s", indexName, fieldName)
}
// Delete field from the index.
if err := index.DeleteField(fieldName); err != nil {
return errors.Wrap(err, "deleting field")

View file

@ -328,7 +328,14 @@ func Test_DBPerShard_GetFieldView2Shards_map_from_RBF(t *testing.T) {
index := "rick"
field := "f"
idx, err := holder.createIndex(index, IndexOptions{})
cim := &CreateIndexMessage{
Index: index,
CreatedAt: 0,
Meta: &IndexOptions{},
}
idx, err := holder.createIndex(cim, false)
panicOn(err)
exp := NewFieldView2Shards()

View file

@ -3585,8 +3585,14 @@ func newTestHolder(tb testing.TB) *Holder {
// fragTestMustOpenIndex returns a new, opened index at a temporary path. Panic on error.
func fragTestMustOpenIndex(index string, holder *Holder, opt IndexOptions) *Index {
cim := &CreateIndexMessage{
Index: index,
CreatedAt: 0,
Meta: &opt,
}
holder.mu.Lock()
idx, err := holder.createIndex(index, opt)
idx, err := holder.createIndex(cim, false)
holder.mu.Unlock()
panicOn(err)

254
holder.go
View file

@ -30,6 +30,7 @@ import (
"syscall"
"time"
"github.com/pilosa/pilosa/v2/disco"
"github.com/pilosa/pilosa/v2/logger"
rbfcfg "github.com/pilosa/pilosa/v2/rbf/cfg"
"github.com/pilosa/pilosa/v2/roaring"
@ -80,6 +81,8 @@ type Holder struct {
opened lockedChan
broadcaster broadcaster
schemator disco.Schemator
serializer Serializer
NewAttrStore func(string) AttrStore
@ -571,7 +574,6 @@ func (h *Holder) Inspect(ctx context.Context, req *InspectRequest) (*HolderInfo,
// Open initializes the root data directory for the holder.
func (h *Holder) Open() error {
h.opening = true
defer func() { h.opening = false }()
@ -854,51 +856,55 @@ func (h *Holder) availableShardsByIndex() map[string]*roaring.Bitmap {
// Schema returns schema information for all indexes, fields, and views.
func (h *Holder) Schema() ([]*IndexInfo, error) {
var a []*IndexInfo
for _, index := range h.Indexes() {
di := &IndexInfo{
Name: index.Name(),
CreatedAt: index.CreatedAt(),
Options: index.Options(),
}
for _, field := range index.Fields() {
fi := &FieldInfo{
Name: field.Name(),
CreatedAt: field.CreatedAt(),
Options: field.Options(),
}
for _, view := range field.views() {
fi.Views = append(fi.Views, &ViewInfo{Name: view.name})
}
sort.Sort(viewInfoSlice(fi.Views))
di.Fields = append(di.Fields, fi)
}
sort.Sort(fieldInfoSlice(di.Fields))
a = append(a, di)
}
sort.Sort(indexInfoSlice(a))
return a, nil
return h.schema(context.TODO(), true)
}
// limitedSchema returns schema information for all indexes and fields.
func (h *Holder) limitedSchema() ([]*IndexInfo, error) {
return h.schema(context.TODO(), false)
}
func (h *Holder) schema(ctx context.Context, includeViews bool) ([]*IndexInfo, error) {
var a []*IndexInfo
for _, index := range h.Indexes() {
di := &IndexInfo{
Name: index.Name(),
CreatedAt: index.CreatedAt(),
Options: index.Options(),
ShardWidth: ShardWidth,
Fields: make([]*FieldInfo, 0, len(index.Fields())),
schema, err := h.schemator.Schema(ctx)
if err != nil {
return nil, errors.Wrapf(err, "getting schema via schemator")
}
for indexName, index := range schema {
cim, err := h.decodeCreateIndexMessage(index.Data)
if err != nil {
return nil, errors.Wrap(err, "decoding CreateIndexMessage")
}
for _, field := range index.Fields() {
if strings.HasPrefix(field.name, "_") {
continue
di := &IndexInfo{
Name: cim.Index,
CreatedAt: cim.CreatedAt,
Options: *cim.Meta,
ShardWidth: ShardWidth,
Fields: make([]*FieldInfo, 0, len(index.Fields)),
}
for fieldName, fieldData := range index.Fields {
createFieldMessage, err := h.decodeCreateFieldMessage(fieldData)
if err != nil {
return nil, errors.Wrap(err, "decoding CreateFieldMessage")
}
fi := &FieldInfo{
Name: field.Name(),
CreatedAt: field.CreatedAt(),
Options: field.Options(),
Name: fieldName,
CreatedAt: createFieldMessage.CreatedAt,
Options: *createFieldMessage.Meta,
}
if includeViews {
// Because views are not stored in etcd, we still rely on the
// local representation of views.
if localField := h.Field(indexName, fieldName); localField != nil {
for _, view := range localField.views() {
fi.Views = append(fi.Views, &ViewInfo{Name: view.name})
}
sort.Sort(viewInfoSlice(fi.Views))
}
}
di.Fields = append(di.Fields, fi)
}
@ -995,7 +1001,66 @@ func (h *Holder) CreateIndex(name string, opt IndexOptions) (*Index, error) {
if h.Index(name) != nil {
return nil, newConflictError(ErrIndexExists)
}
return h.createIndex(name, opt)
cim := &CreateIndexMessage{
Index: name,
CreatedAt: 0,
Meta: &opt,
}
// Create the index in etcd as the system of record.
if err := h.persistIndex(context.Background(), cim); err != nil {
return nil, errors.Wrap(err, "persisting index")
}
return h.createIndex(cim, false)
}
// LoadIndex creates an index based on the information stored in schemator.
// An error is returned if the index already exists.
func (h *Holder) LoadIndex(name string) (*Index, error) {
h.mu.Lock()
defer h.mu.Unlock()
// Ensure index doesn't already exist.
if h.Index(name) != nil {
return nil, newConflictError(ErrIndexExists)
}
return h.loadIndex(name)
}
// LoadField creates a field based on the information stored in schemator.
// An error is returned if the field already exists.
func (h *Holder) LoadField(index, field string) (*Field, error) {
h.mu.Lock()
defer h.mu.Unlock()
// Ensure field doesn't already exist.
if h.Field(index, field) != nil {
return nil, newConflictError(ErrFieldExists)
}
return h.loadField(index, field)
}
// CreateIndexAndBroadcast creates an index locally, then broadcasts the
// creation to other nodes so they can create locally as well. An error is
// returned if the index already exists.
//func (h *Holder) CreateIndexAndBroadcast(name string, opt IndexOptions) (*Index, error) {
func (h *Holder) CreateIndexAndBroadcast(cim *CreateIndexMessage) (*Index, error) {
h.mu.Lock()
defer h.mu.Unlock()
// Ensure index doesn't already exist.
if h.Index(cim.Index) != nil {
return nil, newConflictError(ErrIndexExists)
}
// Create the index in etcd as the system of record.
if err := h.persistIndex(context.Background(), cim); err != nil {
return nil, errors.Wrap(err, "persisting index")
}
return h.createIndex(cim, true)
}
// CreateIndexIfNotExists returns an index by name.
@ -1008,22 +1073,62 @@ func (h *Holder) CreateIndexIfNotExists(name string, opt IndexOptions) (*Index,
if index := h.Index(name); index != nil {
return index, nil
}
return h.createIndex(name, opt)
cim := &CreateIndexMessage{
Index: name,
CreatedAt: 0,
Meta: &opt,
}
// Create the index in etcd as the system of record.
if err := h.persistIndex(context.Background(), cim); err != nil {
// There is a case where the index is not in memory, but it is in
// persistent storage. In that case, this will return an "index exists"
// error, which in that case should return the index. TODO: We may need
// to allow for that in the future.
return nil, errors.Wrap(err, "persisting index")
}
return h.createIndex(cim, false)
}
func (h *Holder) createIndex(name string, opt IndexOptions) (*Index, error) {
if name == "" {
// persistIndex stores the index information in etcd.
func (h *Holder) persistIndex(ctx context.Context, cim *CreateIndexMessage) error {
if cim.Index == "" {
return ErrIndexRequired
}
if err := validateName(cim.Index); err != nil {
return errors.Wrap(err, "validating name")
}
if b, err := h.serializer.Marshal(cim); err != nil {
return errors.Wrap(err, "marshaling")
} else if err := h.schemator.CreateIndex(ctx, cim.Index, b); err != nil {
return errors.Wrapf(err, "writing index to disco: %s", cim.Index)
}
return nil
}
func (h *Holder) createIndex(cim *CreateIndexMessage, broadcast bool) (*Index, error) {
if cim.Index == "" {
return nil, errors.New("index name required")
}
opt := cim.Meta
if opt == nil {
opt = &IndexOptions{}
}
// Otherwise create a new index.
index, err := h.newIndex(h.IndexPath(name), name)
index, err := h.newIndex(h.IndexPath(cim.Index), cim.Index)
if err != nil {
return nil, errors.Wrap(err, "creating")
}
index.keys = opt.Keys
index.trackExistence = opt.TrackExistence
index.createdAt = cim.CreatedAt
if err = index.Open(); err != nil {
return nil, errors.Wrap(err, "opening")
@ -1035,6 +1140,13 @@ func (h *Holder) createIndex(name string, opt IndexOptions) (*Index, error) {
// Update options.
h.addIndex(index)
if broadcast {
// Send the create index message to all nodes.
if err := h.broadcaster.SendSync(cim); err != nil {
return nil, errors.Wrap(err, "sending CreateIndex message")
}
}
// Since this is a new index, we need to kick off
// its translation sync.
if err := h.translationSyncer.Reset(); err != nil {
@ -1044,6 +1156,44 @@ func (h *Holder) createIndex(name string, opt IndexOptions) (*Index, error) {
return index, nil
}
func (h *Holder) loadIndex(indexName string) (*Index, error) {
b, err := h.schemator.Index(context.TODO(), indexName)
if err != nil {
// TODO: we may need to wrap with ConflictError if the error type is
// ErrIndexExists.
return nil, errors.Wrapf(err, "getting index: %s", indexName)
}
cim, err := h.decodeCreateIndexMessage(b)
if err != nil {
return nil, errors.Wrap(err, "decoding CreateIndexMessage")
}
return h.createIndex(cim, false)
}
func (h *Holder) loadField(indexName, fieldName string) (*Field, error) {
b, err := h.schemator.Field(context.TODO(), indexName, fieldName)
if err != nil {
// TODO: we may need to wrap with ConflictError if the error type is
// ErrIndexExists.
return nil, errors.Wrapf(err, "getting field: %s/%s", indexName, fieldName)
}
// Get index.
idx := h.Index(indexName)
if idx == nil {
return nil, errors.Errorf("local index not found: %s", indexName)
}
createFieldMessage, err := h.decodeCreateFieldMessage(b)
if err != nil {
return nil, errors.Wrap(err, "decoding CreateFieldMessage")
}
return idx.createFieldIfNotExists(fieldName, createFieldMessage.Meta)
}
func (h *Holder) newIndex(path, name string) (*Index, error) {
index, err := NewIndex(h, path, name)
if err != nil {
@ -1051,6 +1201,8 @@ func (h *Holder) newIndex(path, name string) (*Index, error) {
}
index.Stats = h.Stats.WithTags(fmt.Sprintf("index:%s", index.Name()))
index.broadcaster = h.broadcaster
index.serializer = h.serializer
index.schemator = h.schemator
index.newAttrStore = h.NewAttrStore
index.columnAttrs = h.NewAttrStore(filepath.Join(index.path, ".data"))
index.OpenTranslateStore = h.OpenTranslateStore
@ -2052,3 +2204,19 @@ func (h *Holder) HasRoaringData() (has bool, err error) {
}
return
}
func (h *Holder) decodeCreateIndexMessage(b []byte) (*CreateIndexMessage, error) {
var cim CreateIndexMessage
if err := h.serializer.Unmarshal(b, &cim); err != nil {
return nil, errors.Wrap(err, "unmarshaling")
}
return &cim, nil
}
func (h *Holder) decodeCreateFieldMessage(b []byte) (*CreateFieldMessage, error) {
var cfm CreateFieldMessage
if err := h.serializer.Unmarshal(b, &cfm); err != nil {
return nil, errors.Wrap(err, "unmarshaling")
}
return &cfm, nil
}

110
index.go
View file

@ -26,6 +26,7 @@ import (
"time"
"github.com/gogo/protobuf/proto"
"github.com/pilosa/pilosa/v2/disco"
"github.com/pilosa/pilosa/v2/internal"
"github.com/pilosa/pilosa/v2/roaring"
"github.com/pilosa/pilosa/v2/stats"
@ -56,6 +57,8 @@ type Index struct {
columnAttrs AttrStore
broadcaster broadcaster
schemator disco.Schemator
serializer Serializer
Stats stats.StatsClient
// Passed to field for foreign-index lookup.
@ -511,7 +514,44 @@ func (i *Index) CreateField(name string, opts ...FieldOption) (*Field, error) {
return nil, errors.Wrap(err, "applying option")
}
return i.createField(name, fo)
cfm := &CreateFieldMessage{
Index: i.name,
Field: name,
CreatedAt: 0,
Meta: fo,
}
// Create the field in etcd as the system of record.
if err := i.persistField(context.Background(), cfm); err != nil {
return nil, errors.Wrap(err, "persisting index")
}
return i.createField(cfm, false)
}
// CreateFieldAndBroadcast creates a field locally, then broadcasts the
// creation to other nodes so they can create locally as well. An error is
// returned if the field already exists.
func (i *Index) CreateFieldAndBroadcast(cfm *CreateFieldMessage) (*Field, error) {
err := validateName(cfm.Field)
if err != nil {
return nil, errors.Wrap(err, "validating name")
}
i.mu.Lock()
defer i.mu.Unlock()
// Ensure field doesn't already exist.
if i.fields[cfm.Field] != nil {
return nil, newConflictError(ErrFieldExists)
}
// Create the field in etcd as the system of record.
if err := i.persistField(context.Background(), cfm); err != nil {
return nil, errors.Wrap(err, "persisting index")
}
return i.createField(cfm, true)
}
// CreateFieldIfNotExists creates a field with the given options if it doesn't exist.
@ -535,7 +575,43 @@ func (i *Index) CreateFieldIfNotExists(name string, opts ...FieldOption) (*Field
return nil, errors.Wrap(err, "applying option")
}
return i.createField(name, fo)
cfm := &CreateFieldMessage{
Index: i.name,
Field: name,
CreatedAt: 0,
Meta: fo,
}
// Create the field in etcd as the system of record.
if err := i.persistField(context.Background(), cfm); err != nil {
// There is a case where the index is not in memory, but it is in
// persistent storage. In that case, this will return an "index exists"
// error, which in that case should return the index. TODO: We may need
// to allow for that in the future.
return nil, errors.Wrap(err, "persisting index")
}
return i.createField(cfm, false)
}
// persistField stores the field information in etcd.
func (i *Index) persistField(ctx context.Context, cfm *CreateFieldMessage) error {
if cfm.Index == "" {
return ErrIndexRequired
} else if cfm.Field == "" {
return ErrFieldRequired
}
if err := validateName(cfm.Field); err != nil {
return errors.Wrap(err, "validating name")
}
if b, err := i.serializer.Marshal(cfm); err != nil {
return errors.Wrap(err, "marshaling")
} else if err := i.schemator.CreateField(ctx, cfm.Index, cfm.Field, b); err != nil {
return errors.Wrapf(err, "writing field to disco: %s/%s", cfm.Index, cfm.Field)
}
return nil
}
func (i *Index) createFieldIfNotExists(name string, opt *FieldOptions) (*Field, error) {
@ -547,21 +623,34 @@ func (i *Index) createFieldIfNotExists(name string, opt *FieldOptions) (*Field,
return f, nil
}
return i.createField(name, opt)
cfm := &CreateFieldMessage{
Index: i.name,
Field: name,
CreatedAt: 0,
Meta: opt,
}
return i.createField(cfm, false)
}
func (i *Index) createField(name string, opt *FieldOptions) (*Field, error) {
if name == "" {
func (i *Index) createField(cfm *CreateFieldMessage, broadcast bool) (*Field, error) {
opt := cfm.Meta
if opt == nil {
opt = &FieldOptions{}
}
if cfm.Field == "" {
return nil, errors.New("field name required")
} else if opt.CacheType != "" && !isValidCacheType(opt.CacheType) {
return nil, ErrInvalidCacheType
}
// Initialize field.
f, err := i.newField(i.fieldPath(name), name)
f, err := i.newField(i.fieldPath(cfm.Field), cfm.Field)
if err != nil {
return nil, errors.Wrap(err, "initializing")
}
f.createdAt = cfm.CreatedAt
// Pass holder through to the field for use in looking
// up a foreign index.
@ -580,11 +669,18 @@ func (i *Index) createField(name string, opt *FieldOptions) (*Field, error) {
}
// Add to index's field lookup.
i.fields[name] = f
i.fields[cfm.Field] = f
// enable Txf to find the index in field_test.go TestField_SetValue
f.idx = i
if broadcast {
// Send the create field message to all nodes.
if err := i.broadcaster.SendSync(cfm); err != nil {
return nil, errors.Wrap(err, "sending CreateField message")
}
}
// Kick off the field's translation sync process.
if err := i.translationSyncer.Reset(); err != nil {
return nil, errors.Wrap(err, "resetting translation syncer")

View file

@ -488,6 +488,8 @@ func NewServer(opts ...ServerOption) (*Server, error) {
s.cluster.confirmDownRetries = s.confirmDownRetries
s.cluster.confirmDownSleep = s.confirmDownSleep
s.holder.broadcaster = s
s.holder.schemator = s.schemator
s.holder.serializer = s.serializer
return s, nil
}
@ -778,14 +780,9 @@ func (s *Server) receiveMessage(m Message) error {
}
case *CreateIndexMessage:
opt := obj.Meta
idx, err := s.holder.CreateIndex(obj.Index, *opt)
if err != nil {
if _, err := s.holder.LoadIndex(obj.Index); err != nil {
return err
}
idx.mu.Lock()
idx.createdAt = obj.CreatedAt
idx.mu.Unlock()
case *DeleteIndexMessage:
if err := s.holder.DeleteIndex(obj.Index); err != nil {
@ -793,18 +790,9 @@ func (s *Server) receiveMessage(m Message) error {
}
case *CreateFieldMessage:
idx := s.holder.Index(obj.Index)
if idx == nil {
return fmt.Errorf("local index not found: %s", obj.Index)
}
opt := obj.Meta
fld, err := idx.createFieldIfNotExists(obj.Field, opt)
if err != nil {
if _, err := s.holder.LoadField(obj.Index, obj.Field); err != nil {
return err
}
fld.mu.Lock()
fld.createdAt = obj.CreatedAt
fld.mu.Unlock()
case *DeleteFieldMessage:
idx := s.holder.Index(obj.Index)

View file

@ -37,7 +37,13 @@ func mustOpenView(tb testing.TB, index, field, name string) *view {
h := NewHolder(path, nil)
// h needs an *Index so we can call h.Index() and get Index.Txf, in TestView_DeleteFragment
idx, err := h.createIndex(index, IndexOptions{})
cim := &CreateIndexMessage{
Index: index,
CreatedAt: 0,
Meta: &IndexOptions{},
}
idx, err := h.createIndex(cim, false)
testhook.Cleanup(tb, func() {
h.Close()
})