This is an automated email from the ASF dual-hosted git repository.

littlecui pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/servicecomb-kie.git


The following commit(s) were added to refs/heads/master by this push:
     new c2eaf67  [feat]get sync value from context instead of directly from 
configuration file (#257)
c2eaf67 is described below

commit c2eaf677dfd29e5548ed549fae17c685eaf78244
Author: kkf1 <[email protected]>
AuthorDate: Sat Aug 20 10:34:56 2022 +0800

    [feat]get sync value from context instead of directly from configuration 
file (#257)
    
    * get sync value from context instead of directly from configuration file
    
    * get sync value from context instead of directly from configuration file
    
    * get sync value from context instead of directly from configuration file
    
    * get sync value from context instead of directly from configuration file
    
    * get sync value from context instead of directly from configuration file
    
    * [feat]get sync value from context instead of directly from configuration 
file
---
 server/datasource/kv_dao_test.go | 45 ++++++++++++++--------------
 server/service/kv/kv_svc.go      | 11 +++----
 server/service/sync/sync.go      | 45 ++++++++++++++++++++++++++++
 server/service/sync/sync_test.go | 65 ++++++++++++++++++++++++++++++++++++++++
 4 files changed, 138 insertions(+), 28 deletions(-)

diff --git a/server/datasource/kv_dao_test.go b/server/datasource/kv_dao_test.go
index 415a432..977944d 100644
--- a/server/datasource/kv_dao_test.go
+++ b/server/datasource/kv_dao_test.go
@@ -25,9 +25,9 @@ import (
 
        "github.com/apache/servicecomb-kie/pkg/common"
        "github.com/apache/servicecomb-kie/pkg/model"
-       "github.com/apache/servicecomb-kie/server/config"
        "github.com/apache/servicecomb-kie/server/datasource"
        kvsvc "github.com/apache/servicecomb-kie/server/service/kv"
+       "github.com/apache/servicecomb-kie/server/service/sync"
        "github.com/apache/servicecomb-kie/test"
        emodel "github.com/apache/servicecomb-service-center/eventbase/model"
        "github.com/apache/servicecomb-service-center/eventbase/service/task"
@@ -112,12 +112,12 @@ func TestWithSync(t *testing.T) {
        if test.IsEmbeddedetcdMode() {
                return
        }
-       // set the sync enabled
-       config.Configurations.Sync.Enabled = true
 
        t.Run("create kv with sync enabled", func(t *testing.T) {
                t.Run("creating a kv will create a task should pass", func(t 
*testing.T) {
-                       kv1, err := kvsvc.Create(context.Background(), 
&model.KVDoc{
+                       // set the sync enabled
+                       ctx := sync.NewContext(context.Background(), true)
+                       kv1, err := kvsvc.Create(ctx, &model.KVDoc{
                                Key:    "sync-create",
                                Value:  "2s",
                                Status: common.StatusEnabled,
@@ -135,31 +135,33 @@ func TestWithSync(t *testing.T) {
                                Domain:  "default",
                                Project: "sync-create",
                        }
-                       _, tempErr := 
kvsvc.FindOneAndDelete(context.Background(), kv1.ID, "sync-create", "default")
+                       _, tempErr := kvsvc.FindOneAndDelete(ctx, kv1.ID, 
"sync-create", "default")
                        assert.Nil(t, tempErr)
-                       resp, tempErr := kvsvc.List(context.Background(), 
"sync-create", "default")
+                       resp, tempErr := kvsvc.List(ctx, "sync-create", 
"default")
                        assert.Nil(t, tempErr)
                        assert.Equal(t, 0, resp.Total)
-                       tasks, tempErr := task.List(context.Background(), 
&listReq)
+                       tasks, tempErr := task.List(ctx, &listReq)
                        assert.Nil(t, tempErr)
                        assert.Equal(t, 2, len(tasks))
-                       tempErr = task.Delete(context.Background(), tasks...)
+                       tempErr = task.Delete(ctx, tasks...)
                        assert.Nil(t, tempErr)
                        tbListReq := emodel.ListTombstoneRequest{
                                Domain:       "default",
                                Project:      "sync-create",
                                ResourceType: datasource.ConfigResource,
                        }
-                       tombstones, tempErr := 
tombstone.List(context.Background(), &tbListReq)
+                       tombstones, tempErr := tombstone.List(ctx, &tbListReq)
                        assert.Equal(t, 1, len(tombstones))
-                       tempErr = tombstone.Delete(context.Background(), 
tombstones...)
+                       tempErr = tombstone.Delete(ctx, tombstones...)
                        assert.Nil(t, tempErr)
                })
        })
 
        t.Run("update kv with sync enabled", func(t *testing.T) {
                t.Run("creating two kvs and updating them will create four 
tasks should pass", func(t *testing.T) {
-                       kv1, err := kvsvc.Create(context.Background(), 
&model.KVDoc{
+                       // set the sync enabled
+                       ctx := sync.NewContext(context.Background(), true)
+                       kv1, err := kvsvc.Create(ctx, &model.KVDoc{
                                Key:    "sync-update-one",
                                Value:  "2s",
                                Status: common.StatusEnabled,
@@ -172,7 +174,7 @@ func TestWithSync(t *testing.T) {
                        })
                        assert.Nil(t, err)
                        assert.NotEmpty(t, kv1.ID)
-                       kv2, err := kvsvc.Create(context.Background(), 
&model.KVDoc{
+                       kv2, err := kvsvc.Create(ctx, &model.KVDoc{
                                Key:    "sync-update-two",
                                Value:  "2s",
                                Status: common.StatusEnabled,
@@ -185,7 +187,7 @@ func TestWithSync(t *testing.T) {
                        })
                        assert.Nil(t, err)
                        assert.NotEmpty(t, kv2.ID)
-                       kv1, tmpErr := kvsvc.Update(context.Background(), 
&model.UpdateKVRequest{
+                       kv1, tmpErr := kvsvc.Update(ctx, &model.UpdateKVRequest{
                                ID:      kv1.ID,
                                Value:   "3s",
                                Domain:  "default",
@@ -193,7 +195,7 @@ func TestWithSync(t *testing.T) {
                        })
                        assert.Nil(t, tmpErr)
                        assert.NotEmpty(t, kv1.ID)
-                       kv2, tmpErr = kvsvc.Update(context.Background(), 
&model.UpdateKVRequest{
+                       kv2, tmpErr = kvsvc.Update(ctx, &model.UpdateKVRequest{
                                ID:      kv2.ID,
                                Value:   "3s",
                                Domain:  "default",
@@ -201,33 +203,30 @@ func TestWithSync(t *testing.T) {
                        })
                        assert.Nil(t, tmpErr)
                        assert.NotEmpty(t, kv2.ID)
-                       _, tempErr := 
kvsvc.FindManyAndDelete(context.Background(), []string{kv1.ID, kv2.ID}, 
"sync-update", "default")
+                       _, tempErr := kvsvc.FindManyAndDelete(ctx, 
[]string{kv1.ID, kv2.ID}, "sync-update", "default")
                        assert.Nil(t, tempErr)
-                       resp, tempErr := kvsvc.List(context.Background(), 
"sync-update", "default")
+                       resp, tempErr := kvsvc.List(ctx, "sync-update", 
"default")
                        assert.Nil(t, tempErr)
                        assert.Equal(t, 0, resp.Total)
                        listReq := emodel.ListTaskRequest{
                                Domain:  "default",
                                Project: "sync-update",
                        }
-                       tasks, tempErr := task.List(context.Background(), 
&listReq)
+                       tasks, tempErr := task.List(ctx, &listReq)
                        assert.Nil(t, tempErr)
                        assert.Equal(t, 6, len(tasks))
-                       tempErr = task.Delete(context.Background(), tasks...)
+                       tempErr = task.Delete(ctx, tasks...)
                        assert.Nil(t, tempErr)
                        tbListReq := emodel.ListTombstoneRequest{
                                Domain:       "default",
                                Project:      "sync-update",
                                ResourceType: datasource.ConfigResource,
                        }
-                       tombstones, tempErr := 
tombstone.List(context.Background(), &tbListReq)
+                       tombstones, tempErr := tombstone.List(ctx, &tbListReq)
                        assert.Equal(t, 2, len(tombstones))
-                       tempErr = tombstone.Delete(context.Background(), 
tombstones...)
+                       tempErr = tombstone.Delete(ctx, tombstones...)
                        assert.Nil(t, tempErr)
                })
        })
 
-       // set the sync unable
-       config.Configurations.Sync.Enabled = false
-
 }
diff --git a/server/service/kv/kv_svc.go b/server/service/kv/kv_svc.go
index ec49ea2..66a958b 100644
--- a/server/service/kv/kv_svc.go
+++ b/server/service/kv/kv_svc.go
@@ -33,9 +33,9 @@ import (
        "github.com/apache/servicecomb-kie/pkg/concurrency"
        "github.com/apache/servicecomb-kie/pkg/model"
        "github.com/apache/servicecomb-kie/pkg/stringutil"
-       cfg "github.com/apache/servicecomb-kie/server/config"
        "github.com/apache/servicecomb-kie/server/datasource"
        "github.com/apache/servicecomb-kie/server/pubsub"
+       "github.com/apache/servicecomb-kie/server/service/sync"
 )
 
 var listSema = concurrency.NewSemaphore(concurrency.DefaultConcurrency)
@@ -113,7 +113,8 @@ func Create(ctx context.Context, kv *model.KVDoc) 
(*model.KVDoc, *errsvc.Error)
                openlog.Error(err.Error())
                return nil, config.NewError(config.ErrInternal, "create kv 
failed")
        }
-       kv, err = datasource.GetBroker().GetKVDao().Create(ctx, kv, 
datasource.WithSync(cfg.GetSync().Enabled))
+
+       kv, err = datasource.GetBroker().GetKVDao().Create(ctx, kv, 
datasource.WithSync(sync.FromContext(ctx)))
        if err != nil {
                openlog.Error(fmt.Sprintf("post err:%s", err.Error()))
                return nil, config.NewError(config.ErrInternal, "create kv 
failed")
@@ -230,7 +231,7 @@ func Update(ctx context.Context, kv *model.UpdateKVRequest) 
(*model.KVDoc, error
        if err != nil {
                return nil, err
        }
-       err = datasource.GetBroker().GetKVDao().Update(ctx, oldKV, 
datasource.WithSync(cfg.GetSync().Enabled))
+       err = datasource.GetBroker().GetKVDao().Update(ctx, oldKV, 
datasource.WithSync(sync.FromContext(ctx)))
        if err != nil {
                return nil, err
        }
@@ -252,7 +253,7 @@ func Update(ctx context.Context, kv *model.UpdateKVRequest) 
(*model.KVDoc, error
 }
 
 func FindOneAndDelete(ctx context.Context, kvID string, project, domain 
string) (*model.KVDoc, error) {
-       kv, err := datasource.GetBroker().GetKVDao().FindOneAndDelete(ctx, 
kvID, project, domain, datasource.WithSync(cfg.GetSync().Enabled))
+       kv, err := datasource.GetBroker().GetKVDao().FindOneAndDelete(ctx, 
kvID, project, domain, datasource.WithSync(sync.FromContext(ctx)))
        if err != nil {
                return nil, err
        }
@@ -272,7 +273,7 @@ func FindManyAndDelete(ctx context.Context, kvIDs []string, 
project, domain stri
        var kvs []*model.KVDoc
        var deleted int64
        var err error
-       kvs, deleted, err = 
datasource.GetBroker().GetKVDao().FindManyAndDelete(ctx, kvIDs, project, 
domain, datasource.WithSync(cfg.GetSync().Enabled))
+       kvs, deleted, err = 
datasource.GetBroker().GetKVDao().FindManyAndDelete(ctx, kvIDs, project, 
domain, datasource.WithSync(sync.FromContext(ctx)))
        if err != nil {
                return nil, err
        }
diff --git a/server/service/sync/sync.go b/server/service/sync/sync.go
new file mode 100644
index 0000000..2320be4
--- /dev/null
+++ b/server/service/sync/sync.go
@@ -0,0 +1,45 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package sync
+
+import (
+       "context"
+)
+
+type ctxKey string
+
+const CtxSyncEnabled ctxKey = "sync"
+
+func NewContext(ctx context.Context, enabled bool) context.Context {
+       if ctx == nil {
+               ctx = context.Background()
+       }
+       return context.WithValue(ctx, CtxSyncEnabled, enabled)
+}
+
+func FromContext(ctx context.Context) bool {
+       if ctx == nil {
+               return false
+       }
+       val := ctx.Value(CtxSyncEnabled)
+       enabled, ok := val.(bool)
+       if !ok {
+               enabled = false
+       }
+       return enabled
+}
diff --git a/server/service/sync/sync_test.go b/server/service/sync/sync_test.go
new file mode 100644
index 0000000..eee695e
--- /dev/null
+++ b/server/service/sync/sync_test.go
@@ -0,0 +1,65 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package sync_test
+
+import (
+       "context"
+       "testing"
+
+       "github.com/apache/servicecomb-kie/server/service/sync"
+       "github.com/stretchr/testify/assert"
+)
+
+func TestNewContext(t *testing.T) {
+       t.Run("parent is nil should new one", func(t *testing.T) {
+               ctx := sync.NewContext(nil, true)
+               assert.NotNil(t, ctx)
+               assert.True(t, sync.FromContext(ctx))
+       })
+       t.Run("parent is not nil should be ok", func(t *testing.T) {
+               ctx := sync.NewContext(context.Background(), true)
+               assert.NotNil(t, ctx)
+               assert.True(t, sync.FromContext(ctx))
+
+               ctx = sync.NewContext(context.Background(), false)
+               assert.False(t, sync.FromContext(ctx))
+       })
+}
+
+func TestFromContext(t *testing.T) {
+       type args struct {
+               ctx context.Context
+       }
+       tests := []struct {
+               name string
+               args args
+               want bool
+       }{
+               {"no ctx should return false", args{nil}, false},
+               {"ctx without sync should return false", 
args{context.Background()}, false},
+               {"ctx with sync=true should return true", 
args{sync.NewContext(context.Background(), true)}, true},
+               {"ctx with sync=false should return true", 
args{sync.NewContext(context.Background(), false)}, false},
+       }
+       for _, tt := range tests {
+               t.Run(tt.name, func(t *testing.T) {
+                       if got := sync.FromContext(tt.args.ctx); got != tt.want 
{
+                               t.Errorf("FromContext() = %v, want %v", got, 
tt.want)
+                       }
+               })
+       }
+}

Reply via email to