mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-11 05:17:54 +00:00
implement ability to update TTL on time fields
This commit is contained in:
parent
e0e01f9f65
commit
0d50bd2890
13 changed files with 1032 additions and 492 deletions
32
api.go
32
api.go
|
|
@ -343,6 +343,38 @@ func (api *API) CreateField(ctx context.Context, indexName string, fieldName str
|
|||
return field, nil
|
||||
}
|
||||
|
||||
// FieldUpdate represents a change to a field. The thinking is to only
|
||||
// support changing one field option at a time to keep the
|
||||
// implementation sane. At time of writing, only TTL is supported.
|
||||
type FieldUpdate struct {
|
||||
Option string `json:"option"`
|
||||
Value string `json:"value"`
|
||||
}
|
||||
|
||||
func (api *API) UpdateField(ctx context.Context, indexName, fieldName string, update FieldUpdate) error {
|
||||
// Find index.
|
||||
index := api.holder.Index(indexName)
|
||||
if index == nil {
|
||||
return newNotFoundError(ErrIndexNotFound, indexName)
|
||||
}
|
||||
|
||||
cfm, err := index.UpdateField(ctx, fieldName, update)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "updating field")
|
||||
}
|
||||
|
||||
if err := index.UpdateFieldLocal(cfm, update); err != nil {
|
||||
return errors.Wrap(err, "updating field locally")
|
||||
}
|
||||
|
||||
// broadcast field update
|
||||
err = api.holder.sendOrSpool(&UpdateFieldMessage{
|
||||
CreateFieldMessage: *cfm,
|
||||
Update: update,
|
||||
})
|
||||
return errors.Wrap(err, "sending UpdateField message")
|
||||
}
|
||||
|
||||
// Field retrieves the named field.
|
||||
func (api *API) Field(ctx context.Context, indexName, fieldName string) (*Field, error) {
|
||||
span, _ := tracing.StartSpanFromContext(ctx, "API.Field")
|
||||
|
|
|
|||
|
|
@ -70,6 +70,7 @@ const (
|
|||
messageTypeTransaction
|
||||
messageTypeResizeNodeMessage
|
||||
messageTypeResizeAbortMessage
|
||||
messageTypeUpdateField
|
||||
)
|
||||
|
||||
// MarshalInternalMessage serializes the pilosa message and adds pilosa internal
|
||||
|
|
@ -121,6 +122,8 @@ func getMessage(typ byte) Message {
|
|||
return &ResizeNodeMessage{}
|
||||
case messageTypeResizeAbortMessage:
|
||||
return &ResizeAbortMessage{}
|
||||
case messageTypeUpdateField:
|
||||
return &UpdateFieldMessage{}
|
||||
default:
|
||||
panic(fmt.Sprintf("unknown message type %d", typ))
|
||||
}
|
||||
|
|
@ -164,6 +167,8 @@ func getMessageType(m Message) byte {
|
|||
return messageTypeResizeNodeMessage
|
||||
case *ResizeAbortMessage:
|
||||
return messageTypeResizeAbortMessage
|
||||
case *UpdateFieldMessage:
|
||||
return messageTypeUpdateField
|
||||
default:
|
||||
panic(fmt.Sprintf("don't have type for message %#v", m))
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1991,6 +1991,14 @@ type CreateFieldMessage struct {
|
|||
Meta *FieldOptions
|
||||
}
|
||||
|
||||
// UpdateFieldMessage represents a change to an existing field. The
|
||||
// CreateFieldMessage holds the changed field, while the update shows
|
||||
// the change that was made.
|
||||
type UpdateFieldMessage struct {
|
||||
CreateFieldMessage CreateFieldMessage
|
||||
Update FieldUpdate
|
||||
}
|
||||
|
||||
// DeleteFieldMessage is an internal message indicating field deletion.
|
||||
type DeleteFieldMessage struct {
|
||||
Index string
|
||||
|
|
|
|||
|
|
@ -127,6 +127,7 @@ type Schemator interface {
|
|||
DeleteIndex(ctx context.Context, name string) error
|
||||
Field(ctx context.Context, index, field string) ([]byte, error)
|
||||
CreateField(ctx context.Context, index, field string, val []byte) error
|
||||
UpdateField(ctx context.Context, index, field string, val []byte) error
|
||||
DeleteField(ctx context.Context, index, field string) error
|
||||
View(ctx context.Context, index, field, view string) (bool, error)
|
||||
CreateView(ctx context.Context, index, field, view string) error
|
||||
|
|
@ -285,6 +286,11 @@ func (*nopSchemator) CreateField(ctx context.Context, index, field string, val [
|
|||
return nil
|
||||
}
|
||||
|
||||
// UpdateField is a no-op implementation of the Schemator UpdateField method.
|
||||
func (*nopSchemator) UpdateField(ctx context.Context, index, field string, val []byte) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// DeleteField is a no-op implementation of the Schemator DeleteField method.
|
||||
func (*nopSchemator) DeleteField(ctx context.Context, index, field string) error { return nil }
|
||||
|
||||
|
|
@ -401,6 +407,24 @@ func (s *inMemSchemator) CreateField(ctx context.Context, index, field string, v
|
|||
return nil
|
||||
}
|
||||
|
||||
func (s *inMemSchemator) UpdateField(ctx context.Context, index, field string, val []byte) error {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
idx, ok := s.schema[index]
|
||||
if !ok {
|
||||
return ErrIndexDoesNotExist
|
||||
}
|
||||
if fld, ok := idx.Fields[field]; ok {
|
||||
// The current logic in pilosa doesn't allow us to return ErrFieldExists
|
||||
// here, so for now we just update the Data value if the field already
|
||||
// exists.
|
||||
fld.Data = val
|
||||
return nil
|
||||
} else {
|
||||
return ErrFieldDoesNotExist
|
||||
}
|
||||
}
|
||||
|
||||
// DeleteField is an in-memory implementation of the Schemator DeleteField method.
|
||||
func (s *inMemSchemator) DeleteField(ctx context.Context, index, field string) error {
|
||||
s.mu.Lock()
|
||||
|
|
|
|||
|
|
@ -70,6 +70,14 @@ func (s Serializer) Unmarshal(buf []byte, m pilosa.Message) error {
|
|||
}
|
||||
s.decodeCreateFieldMessage(msg, mt)
|
||||
return nil
|
||||
case *pilosa.UpdateFieldMessage:
|
||||
msg := &pb.UpdateFieldMessage{}
|
||||
err := proto.Unmarshal(buf, msg)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "unmarshaling UpdateFieldMessage")
|
||||
}
|
||||
s.decodeUpdateFieldMessage(msg, mt)
|
||||
return nil
|
||||
case *pilosa.DeleteFieldMessage:
|
||||
msg := &pb.DeleteFieldMessage{}
|
||||
err := proto.Unmarshal(buf, msg)
|
||||
|
|
@ -339,6 +347,8 @@ func (s Serializer) encodeToProto(m pilosa.Message) proto.Message {
|
|||
return s.encodeDeleteIndexMessage(mt)
|
||||
case *pilosa.CreateFieldMessage:
|
||||
return s.encodeCreateFieldMessage(mt)
|
||||
case *pilosa.UpdateFieldMessage:
|
||||
return s.encodeUpdateFieldMessage(mt)
|
||||
case *pilosa.DeleteFieldMessage:
|
||||
return s.encodeDeleteFieldMessage(mt)
|
||||
case *pilosa.DeleteAvailableShardMessage:
|
||||
|
|
@ -751,6 +761,20 @@ func (s Serializer) encodeCreateFieldMessage(m *pilosa.CreateFieldMessage) *pb.C
|
|||
}
|
||||
}
|
||||
|
||||
func (s Serializer) encodeUpdateFieldMessage(m *pilosa.UpdateFieldMessage) *pb.UpdateFieldMessage {
|
||||
return &pb.UpdateFieldMessage{
|
||||
CreateFieldMessage: s.encodeCreateFieldMessage(&m.CreateFieldMessage),
|
||||
Update: s.encodeFieldUpdate(&m.Update),
|
||||
}
|
||||
}
|
||||
|
||||
func (s Serializer) encodeFieldUpdate(m *pilosa.FieldUpdate) *pb.FieldUpdate {
|
||||
return &pb.FieldUpdate{
|
||||
Option: m.Option,
|
||||
Value: m.Value,
|
||||
}
|
||||
}
|
||||
|
||||
func (s Serializer) encodeDeleteFieldMessage(m *pilosa.DeleteFieldMessage) *pb.DeleteFieldMessage {
|
||||
return &pb.DeleteFieldMessage{
|
||||
Index: m.Index,
|
||||
|
|
@ -1149,6 +1173,16 @@ func (s Serializer) decodeCreateFieldMessage(pb *pb.CreateFieldMessage, m *pilos
|
|||
s.decodeFieldOptions(pb.Meta, m.Meta)
|
||||
}
|
||||
|
||||
func (s Serializer) decodeUpdateFieldMessage(pb *pb.UpdateFieldMessage, m *pilosa.UpdateFieldMessage) {
|
||||
s.decodeCreateFieldMessage(pb.CreateFieldMessage, &m.CreateFieldMessage)
|
||||
s.decodeFieldUpdate(pb.Update, &m.Update)
|
||||
}
|
||||
|
||||
func (s Serializer) decodeFieldUpdate(pb *pb.FieldUpdate, m *pilosa.FieldUpdate) {
|
||||
m.Option = pb.Option
|
||||
m.Value = pb.Value
|
||||
}
|
||||
|
||||
func (s Serializer) decodeDeleteFieldMessage(pb *pb.DeleteFieldMessage, m *pilosa.DeleteFieldMessage) {
|
||||
m.Index = pb.Index
|
||||
m.Field = pb.Field
|
||||
|
|
|
|||
|
|
@ -906,6 +906,33 @@ func (e *Etcd) CreateField(ctx context.Context, indexName string, name string, v
|
|||
return nil
|
||||
}
|
||||
|
||||
func (e *Etcd) UpdateField(ctx context.Context, indexName string, name string, val []byte) error {
|
||||
key := schemaPrefix + indexName + "/" + name
|
||||
|
||||
// Set up Op to write field value as bytes.
|
||||
op := clientv3.OpPut(key, "")
|
||||
op.WithValueBytes(val)
|
||||
|
||||
// Check for key existence, and execute Op within a transaction.
|
||||
var resp *clientv3.TxnResponse
|
||||
|
||||
err := e.retryClient(func(cli *clientv3.Client) (err error) {
|
||||
resp, err = cli.Txn(ctx).
|
||||
If(clientv3util.KeyExists(key)).
|
||||
Then(op).
|
||||
Commit()
|
||||
return err
|
||||
})
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "executing transaction")
|
||||
}
|
||||
|
||||
if !resp.Succeeded {
|
||||
return disco.ErrFieldDoesNotExist
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (e *Etcd) DeleteField(ctx context.Context, indexname string, name string) (err error) {
|
||||
key := schemaPrefix + indexname + "/" + name
|
||||
// Deleting field and views in one transaction.
|
||||
|
|
|
|||
|
|
@ -432,6 +432,7 @@ func newRouter(handler *Handler) http.Handler {
|
|||
router.HandleFunc("/index/{index}/field/{field}/view", handler.chkAuthZ(handler.handleGetView, authz.Admin)).Methods("GET")
|
||||
router.HandleFunc("/index/{index}/field/{field}/view/{view}", handler.chkAuthZ(handler.handleDeleteView, authz.Admin)).Methods("DELETE").Name("DeleteView")
|
||||
router.HandleFunc("/index/{index}/field/{field}", handler.chkAuthZ(handler.handlePostField, authz.Write)).Methods("POST").Name("PostField")
|
||||
router.HandleFunc("/index/{index}/field/{field}", handler.chkAuthZ(handler.handlePatchField, authz.Write)).Methods("PATCH").Name("PatchField")
|
||||
router.HandleFunc("/index/{index}/field/{field}", handler.chkAuthZ(handler.handleDeleteField, authz.Write)).Methods("DELETE").Name("DeleteField")
|
||||
router.HandleFunc("/index/{index}/field/{field}/import", handler.chkAuthZ(handler.handlePostImport, authz.Write)).Methods("POST").Name("PostImport")
|
||||
router.HandleFunc("/index/{index}/field/{field}/mutex-check", handler.chkAuthZ(handler.handleGetMutexCheck, authz.Read)).Methods("GET").Name("GetMutexCheck")
|
||||
|
|
@ -794,6 +795,7 @@ func (r *successResponse) check(err error) (statusCode int) {
|
|||
case ConflictError:
|
||||
statusCode = http.StatusConflict
|
||||
case NotFoundError:
|
||||
// TODO I think any error matches NotFoundError because it's a `type NotFoundError error`
|
||||
statusCode = http.StatusNotFound
|
||||
default:
|
||||
statusCode = http.StatusInternalServerError
|
||||
|
|
@ -1699,6 +1701,41 @@ func (h *Handler) handlePostField(w http.ResponseWriter, r *http.Request) {
|
|||
resp.write(w, err)
|
||||
}
|
||||
|
||||
// handlePatchField handles updates to field schema at /index/{index}/field/{field}
|
||||
func (h *Handler) handlePatchField(w http.ResponseWriter, r *http.Request) {
|
||||
if !validHeaderAcceptJSON(r.Header) {
|
||||
http.Error(w, "JSON only acceptable response", http.StatusNotAcceptable)
|
||||
return
|
||||
}
|
||||
|
||||
indexName, ok := mux.Vars(r)["index"]
|
||||
if !ok {
|
||||
http.Error(w, "index name is required", http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
|
||||
fieldName, ok := mux.Vars(r)["field"]
|
||||
if !ok {
|
||||
http.Error(w, "field name is required", http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
|
||||
resp := successResponse{h: h, Name: fieldName}
|
||||
|
||||
// Decode request.
|
||||
var req FieldUpdate
|
||||
dec := json.NewDecoder(r.Body)
|
||||
dec.DisallowUnknownFields()
|
||||
err := dec.Decode(&req)
|
||||
if err != nil && err != io.EOF {
|
||||
resp.write(w, err)
|
||||
return
|
||||
}
|
||||
|
||||
err = h.api.UpdateField(r.Context(), indexName, fieldName, req)
|
||||
resp.write(w, err)
|
||||
}
|
||||
|
||||
// handlePostIngestData handles JSON ingest data that may need key
|
||||
// translation, for the entire cluster.
|
||||
func (h *Handler) handlePostIngestData(w http.ResponseWriter, r *http.Request) {
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@
|
|||
package pilosa_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"net"
|
||||
|
|
@ -10,6 +11,7 @@ import (
|
|||
"sort"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
pilosa "github.com/molecula/featurebase/v3"
|
||||
"github.com/molecula/featurebase/v3/encoding/proto"
|
||||
|
|
@ -80,6 +82,38 @@ func TestMarshalUnmarshalTransactionResponse(t *testing.T) {
|
|||
}
|
||||
}
|
||||
|
||||
func TestUpdateFieldTTL(t *testing.T) {
|
||||
c := test.MustRunCluster(t, 3)
|
||||
defer c.Close()
|
||||
|
||||
c.CreateField(t, "ttltest", pilosa.IndexOptions{}, "timefield", pilosa.OptFieldTypeTime(pilosa.TimeQuantum("YMD"), "0"))
|
||||
|
||||
nodeURL := c.Nodes[0].URL() + "/index/ttltest/field/timefield"
|
||||
fmt.Println(nodeURL)
|
||||
|
||||
req, err := gohttp.NewRequest("PATCH", nodeURL, strings.NewReader(`{"option": "ttl", "value": "48h"}`))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
resp, err := gohttp.DefaultClient.Do(req)
|
||||
if err != nil {
|
||||
t.Fatalf("doing option request: %v", err)
|
||||
} else if resp.StatusCode != 200 {
|
||||
t.Fatalf("unexpected status updating TTL")
|
||||
}
|
||||
|
||||
for _, node := range c.Nodes {
|
||||
ii, err := node.API.Schema(context.Background(), false)
|
||||
if err != nil {
|
||||
t.Fatalf("getting schema: %v", err)
|
||||
}
|
||||
if ii[0].Fields[0].Options.TTL != time.Hour*48 {
|
||||
t.Fatalf("unexpected TTL after update: %s", ii[0].Fields[0].Options.TTL)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestIngestSchemaHandler(t *testing.T) {
|
||||
c := test.MustRunCluster(t, 3,
|
||||
[]server.CommandOption{
|
||||
|
|
|
|||
73
index.go
73
index.go
|
|
@ -9,6 +9,7 @@ import (
|
|||
"sort"
|
||||
"strconv"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/molecula/featurebase/v3/disco"
|
||||
"github.com/molecula/featurebase/v3/pql"
|
||||
|
|
@ -719,6 +720,78 @@ func (i *Index) persistField(ctx context.Context, cfm *CreateFieldMessage) error
|
|||
return nil
|
||||
}
|
||||
|
||||
func (i *Index) persistUpdateField(ctx context.Context, cfm *CreateFieldMessage) error {
|
||||
if cfm.Index == "" {
|
||||
return ErrIndexRequired
|
||||
} else if cfm.Field == "" {
|
||||
return ErrFieldRequired
|
||||
}
|
||||
|
||||
if b, err := i.serializer.Marshal(cfm); err != nil {
|
||||
return errors.Wrap(err, "marshaling")
|
||||
} else if err := i.Schemator.UpdateField(ctx, cfm.Index, cfm.Field, b); errors.Cause(err) == disco.ErrFieldDoesNotExist {
|
||||
return ErrFieldNotFound
|
||||
} else if err != nil {
|
||||
return errors.Wrapf(err, "writing field to disco: %s/%s", cfm.Index, cfm.Field)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (i *Index) UpdateField(ctx context.Context, name string, update FieldUpdate) (*CreateFieldMessage, error) {
|
||||
// Get field from etcd
|
||||
buf, err := i.Schemator.Field(ctx, i.name, name)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "getting field '%s' from etcd", name)
|
||||
}
|
||||
cfm, err := decodeCreateFieldMessage(i.holder.serializer, buf)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "decoding CreateFieldMessage")
|
||||
} else if cfm == nil {
|
||||
return nil, errors.New("got nil CreateFieldMessage when decoding")
|
||||
}
|
||||
|
||||
// Handle the options we know how to update, or error.
|
||||
switch update.Option {
|
||||
case "TTL", "ttl":
|
||||
if cfm.Meta.Type != FieldTypeTime {
|
||||
return nil, NewBadRequestError(errors.Errorf("can only add TTL to a 'time' type field, not '%s'", cfm.Meta.Type))
|
||||
}
|
||||
dur, err := time.ParseDuration(update.Value)
|
||||
if err != nil {
|
||||
return nil, NewBadRequestError(errors.Wrap(err, "parsing duration"))
|
||||
}
|
||||
cfm.Meta.TTL = dur
|
||||
default:
|
||||
return nil, NewBadRequestError(errors.Errorf("updates for option '%s' are not supported", update.Option))
|
||||
}
|
||||
|
||||
// Persist the updated field to etcd.
|
||||
if err := i.persistUpdateField(ctx, cfm); err != nil {
|
||||
return nil, errors.Wrap(err, "persisting updated field")
|
||||
}
|
||||
|
||||
return cfm, nil
|
||||
}
|
||||
|
||||
func (i *Index) UpdateFieldLocal(cfm *CreateFieldMessage, update FieldUpdate) error {
|
||||
// Update local structures. This assumes we don't need to do
|
||||
// anything else... which is fine for TTL specifically, but I'm
|
||||
// not sure about other things, so be aware when adding new update
|
||||
// abilities.
|
||||
i.mu.Lock()
|
||||
defer i.mu.Unlock()
|
||||
field := i.field(cfm.Field)
|
||||
if field == nil {
|
||||
return errors.Errorf("field '%s' not found locally", cfm.Field)
|
||||
}
|
||||
if err := field.applyOptions(*cfm.Meta); err != nil {
|
||||
return errors.Wrap(err, "updating local field options")
|
||||
}
|
||||
|
||||
return nil
|
||||
|
||||
}
|
||||
|
||||
// createFieldIfNotExists creates the field if it does not already exist in the
|
||||
// in-memory index structure. This is not related to whether or not the field
|
||||
// exists in etcd.
|
||||
|
|
|
|||
1060
pb/private.pb.go
1060
pb/private.pb.go
File diff suppressed because it is too large
Load diff
|
|
@ -76,6 +76,16 @@ message CreateFieldMessage {
|
|||
int64 CreatedAt = 4;
|
||||
}
|
||||
|
||||
message UpdateFieldMessage {
|
||||
CreateFieldMessage CreateFieldMessage = 1;
|
||||
FieldUpdate Update = 2;
|
||||
}
|
||||
|
||||
message FieldUpdate {
|
||||
string Option = 1;
|
||||
string Value = 2;
|
||||
}
|
||||
|
||||
message DeleteFieldMessage {
|
||||
string Index = 1;
|
||||
string Field = 2;
|
||||
|
|
|
|||
175
pb/public.pb.go
175
pb/public.pb.go
|
|
@ -6149,10 +6149,7 @@ func (m *Row) Unmarshal(dAtA []byte) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if skippy < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) < 0 {
|
||||
if (skippy < 0) || (iNdEx+skippy) < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) > l {
|
||||
|
|
@ -6237,10 +6234,7 @@ func (m *RowMatrix) Unmarshal(dAtA []byte) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if skippy < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) < 0 {
|
||||
if (skippy < 0) || (iNdEx+skippy) < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) > l {
|
||||
|
|
@ -6363,10 +6357,7 @@ func (m *SignedRow) Unmarshal(dAtA []byte) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if skippy < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) < 0 {
|
||||
if (skippy < 0) || (iNdEx+skippy) < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) > l {
|
||||
|
|
@ -6525,10 +6516,7 @@ func (m *RowIdentifiers) Unmarshal(dAtA []byte) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if skippy < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) < 0 {
|
||||
if (skippy < 0) || (iNdEx+skippy) < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) > l {
|
||||
|
|
@ -6655,10 +6643,7 @@ func (m *IDList) Unmarshal(dAtA []byte) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if skippy < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) < 0 {
|
||||
if (skippy < 0) || (iNdEx+skippy) < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) > l {
|
||||
|
|
@ -6762,10 +6747,7 @@ func (m *ExtractedIDColumn) Unmarshal(dAtA []byte) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if skippy < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) < 0 {
|
||||
if (skippy < 0) || (iNdEx+skippy) < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) > l {
|
||||
|
|
@ -6882,10 +6864,7 @@ func (m *ExtractedIDMatrix) Unmarshal(dAtA []byte) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if skippy < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) < 0 {
|
||||
if (skippy < 0) || (iNdEx+skippy) < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) > l {
|
||||
|
|
@ -6968,10 +6947,7 @@ func (m *KeyList) Unmarshal(dAtA []byte) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if skippy < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) < 0 {
|
||||
if (skippy < 0) || (iNdEx+skippy) < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) > l {
|
||||
|
|
@ -7185,10 +7161,7 @@ func (m *ExtractedTableValue) Unmarshal(dAtA []byte) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if skippy < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) < 0 {
|
||||
if (skippy < 0) || (iNdEx+skippy) < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) > l {
|
||||
|
|
@ -7325,10 +7298,7 @@ func (m *ExtractedTableColumn) Unmarshal(dAtA []byte) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if skippy < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) < 0 {
|
||||
if (skippy < 0) || (iNdEx+skippy) < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) > l {
|
||||
|
|
@ -7443,10 +7413,7 @@ func (m *ExtractedTableField) Unmarshal(dAtA []byte) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if skippy < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) < 0 {
|
||||
if (skippy < 0) || (iNdEx+skippy) < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) > l {
|
||||
|
|
@ -7565,10 +7532,7 @@ func (m *ExtractedTable) Unmarshal(dAtA []byte) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if skippy < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) < 0 {
|
||||
if (skippy < 0) || (iNdEx+skippy) < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) > l {
|
||||
|
|
@ -7689,10 +7653,7 @@ func (m *Pair) Unmarshal(dAtA []byte) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if skippy < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) < 0 {
|
||||
if (skippy < 0) || (iNdEx+skippy) < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) > l {
|
||||
|
|
@ -7811,10 +7772,7 @@ func (m *PairField) Unmarshal(dAtA []byte) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if skippy < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) < 0 {
|
||||
if (skippy < 0) || (iNdEx+skippy) < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) > l {
|
||||
|
|
@ -7931,10 +7889,7 @@ func (m *PairsField) Unmarshal(dAtA []byte) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if skippy < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) < 0 {
|
||||
if (skippy < 0) || (iNdEx+skippy) < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) > l {
|
||||
|
|
@ -8004,10 +7959,7 @@ func (m *Int64) Unmarshal(dAtA []byte) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if skippy < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) < 0 {
|
||||
if (skippy < 0) || (iNdEx+skippy) < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) > l {
|
||||
|
|
@ -8177,10 +8129,7 @@ func (m *FieldRow) Unmarshal(dAtA []byte) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if skippy < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) < 0 {
|
||||
if (skippy < 0) || (iNdEx+skippy) < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) > l {
|
||||
|
|
@ -8303,10 +8252,7 @@ func (m *GroupCount) Unmarshal(dAtA []byte) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if skippy < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) < 0 {
|
||||
if (skippy < 0) || (iNdEx+skippy) < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) > l {
|
||||
|
|
@ -8474,10 +8420,7 @@ func (m *ValCount) Unmarshal(dAtA []byte) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if skippy < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) < 0 {
|
||||
if (skippy < 0) || (iNdEx+skippy) < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) > l {
|
||||
|
|
@ -8566,10 +8509,7 @@ func (m *Decimal) Unmarshal(dAtA []byte) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if skippy < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) < 0 {
|
||||
if (skippy < 0) || (iNdEx+skippy) < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) > l {
|
||||
|
|
@ -8684,10 +8624,7 @@ func (m *DistinctTimestamp) Unmarshal(dAtA []byte) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if skippy < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) < 0 {
|
||||
if (skippy < 0) || (iNdEx+skippy) < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) > l {
|
||||
|
|
@ -8939,10 +8876,7 @@ func (m *QueryRequest) Unmarshal(dAtA []byte) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if skippy < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) < 0 {
|
||||
if (skippy < 0) || (iNdEx+skippy) < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) > l {
|
||||
|
|
@ -9059,10 +8993,7 @@ func (m *QueryResponse) Unmarshal(dAtA []byte) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if skippy < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) < 0 {
|
||||
if (skippy < 0) || (iNdEx+skippy) < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) > l {
|
||||
|
|
@ -9711,10 +9642,7 @@ func (m *QueryResult) Unmarshal(dAtA []byte) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if skippy < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) < 0 {
|
||||
if (skippy < 0) || (iNdEx+skippy) < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) > l {
|
||||
|
|
@ -10198,10 +10126,7 @@ func (m *ImportRequest) Unmarshal(dAtA []byte) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if skippy < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) < 0 {
|
||||
if (skippy < 0) || (iNdEx+skippy) < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) > l {
|
||||
|
|
@ -10663,10 +10588,7 @@ func (m *ImportValueRequest) Unmarshal(dAtA []byte) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if skippy < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) < 0 {
|
||||
if (skippy < 0) || (iNdEx+skippy) < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) > l {
|
||||
|
|
@ -10836,10 +10758,7 @@ func (m *AtomicRecord) Unmarshal(dAtA []byte) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if skippy < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) < 0 {
|
||||
if (skippy < 0) || (iNdEx+skippy) < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) > l {
|
||||
|
|
@ -10922,10 +10841,7 @@ func (m *AtomicImportResponse) Unmarshal(dAtA []byte) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if skippy < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) < 0 {
|
||||
if (skippy < 0) || (iNdEx+skippy) < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) > l {
|
||||
|
|
@ -11092,10 +11008,7 @@ func (m *TranslateKeysRequest) Unmarshal(dAtA []byte) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if skippy < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) < 0 {
|
||||
if (skippy < 0) || (iNdEx+skippy) < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) > l {
|
||||
|
|
@ -11222,10 +11135,7 @@ func (m *TranslateKeysResponse) Unmarshal(dAtA []byte) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if skippy < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) < 0 {
|
||||
if (skippy < 0) || (iNdEx+skippy) < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) > l {
|
||||
|
|
@ -11416,10 +11326,7 @@ func (m *TranslateIDsRequest) Unmarshal(dAtA []byte) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if skippy < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) < 0 {
|
||||
if (skippy < 0) || (iNdEx+skippy) < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) > l {
|
||||
|
|
@ -11502,10 +11409,7 @@ func (m *TranslateIDsResponse) Unmarshal(dAtA []byte) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if skippy < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) < 0 {
|
||||
if (skippy < 0) || (iNdEx+skippy) < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) > l {
|
||||
|
|
@ -11622,10 +11526,7 @@ func (m *ImportRoaringRequestView) Unmarshal(dAtA []byte) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if skippy < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) < 0 {
|
||||
if (skippy < 0) || (iNdEx+skippy) < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) > l {
|
||||
|
|
@ -11839,10 +11740,7 @@ func (m *ImportRoaringRequest) Unmarshal(dAtA []byte) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if skippy < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) < 0 {
|
||||
if (skippy < 0) || (iNdEx+skippy) < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) > l {
|
||||
|
|
@ -11959,10 +11857,7 @@ func (m *GroupCounts) Unmarshal(dAtA []byte) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if skippy < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) < 0 {
|
||||
if (skippy < 0) || (iNdEx+skippy) < 0 {
|
||||
return ErrInvalidLengthPublic
|
||||
}
|
||||
if (iNdEx + skippy) > l {
|
||||
|
|
|
|||
|
|
@ -990,6 +990,11 @@ func (s *Server) receiveMessage(m Message) error {
|
|||
return err
|
||||
}
|
||||
|
||||
case *UpdateFieldMessage:
|
||||
idx := s.holder.Index(obj.CreateFieldMessage.Index)
|
||||
if err := idx.UpdateFieldLocal(&obj.CreateFieldMessage, obj.Update); err != nil {
|
||||
return err
|
||||
}
|
||||
case *DeleteFieldMessage:
|
||||
idx := s.holder.Index(obj.Index)
|
||||
if err := idx.DeleteField(obj.Field); err != nil {
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue