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

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


The following commit(s) were added to refs/heads/master by this push:
     new 157f175  [feat] add service in eventbase
     new 197d96a  Merge pull request #1182 from robotLJW/master
157f175 is described below

commit 157f175501064e50eabbf64558e34671272357da
Author: robotljw <[email protected]>
AuthorDate: Tue Dec 21 22:17:24 2021 +0800

    [feat] add service in eventbase
---
 .../bootstrap.go}                                  | 17 +++++----
 eventbase/datasource/etcd/etcd.go                  | 10 ++---
 eventbase/datasource/etcd/etcd_test.go             |  8 ++--
 eventbase/datasource/etcd/key/key.go               |  6 +++
 eventbase/datasource/etcd/task/task_dao.go         |  8 ++--
 eventbase/datasource/etcd/task/task_dao_test.go    | 17 ++++++---
 .../datasource/etcd/tombstone/tombstone_dao.go     | 12 +++---
 .../etcd/tombstone/tombstone_dao_test.go           | 19 ++++++----
 eventbase/datasource/manager.go                    |  1 +
 eventbase/datasource/mongo/client/client.go        |  4 +-
 eventbase/datasource/mongo/{ => model}/types.go    |  2 +-
 eventbase/datasource/mongo/mongo.go                | 33 ++++++++--------
 eventbase/datasource/mongo/task/task_dao.go        | 44 ++++++++++++----------
 eventbase/datasource/mongo/task/task_dao_test.go   | 14 ++++---
 .../datasource/mongo/tombstone/tombstone_dao.go    | 43 +++++++++++----------
 .../mongo/tombstone/tombstone_dao_test.go          | 14 ++++---
 eventbase/datasource/options.go                    | 32 ++++++++++++++++
 eventbase/datasource/task.go                       |  2 +-
 eventbase/datasource/tlsutil/tlsutil_test.go       |  2 +-
 eventbase/datasource/tombstone.go                  |  6 +--
 eventbase/go.mod                                   |  2 +-
 eventbase/{model => request}/tombstone_request.go  | 19 +++++++++-
 .../tombstone.go => service/task/task_svc.go}      | 33 +++++++++++-----
 .../tombstone/tombstone_svc.go}                    | 28 +++++++++-----
 24 files changed, 242 insertions(+), 134 deletions(-)

diff --git a/eventbase/model/tombstone_request.go 
b/eventbase/bootstrap/bootstrap.go
similarity index 65%
copy from eventbase/model/tombstone_request.go
copy to eventbase/bootstrap/bootstrap.go
index 1e77f62..5613a80 100644
--- a/eventbase/model/tombstone_request.go
+++ b/eventbase/bootstrap/bootstrap.go
@@ -15,12 +15,13 @@
  * limitations under the License.
  */
 
-package model
+package bootstrap
 
-// GetTombstoneRequest contains tombstone get request params
-type GetTombstoneRequest struct {
-       Project      string `json:"project,omitempty" yaml:"project,omitempty"`
-       Domain       string `json:"domain,omitempty" yaml:"domain,omitempty"`
-       ResourceType string `json:"resource_type,omitempty" 
yaml:"resource_type,omitempty"`
-       ResourceID   string `json:"resource_id,omitempty" 
yaml:"resource_id,omitempty"`
-}
+import (
+       // support embedded etcd
+       _ "github.com/little-cui/etcdadpt/embedded"
+       _ "github.com/little-cui/etcdadpt/remote"
+
+       _ 
"github.com/apache/servicecomb-service-center/eventbase/datasource/etcd"
+       _ 
"github.com/apache/servicecomb-service-center/eventbase/datasource/mongo"
+)
diff --git a/eventbase/datasource/etcd/etcd.go 
b/eventbase/datasource/etcd/etcd.go
index a1ca285..a5389cb 100644
--- a/eventbase/datasource/etcd/etcd.go
+++ b/eventbase/datasource/etcd/etcd.go
@@ -24,13 +24,11 @@ import (
        "github.com/go-chassis/cari/db"
        "github.com/go-chassis/openlog"
        "github.com/little-cui/etcdadpt"
-       _ "github.com/little-cui/etcdadpt/embedded"
-       _ "github.com/little-cui/etcdadpt/remote"
 
-       "servicecomb-service-center/eventbase/datasource"
-       "servicecomb-service-center/eventbase/datasource/etcd/task"
-       "servicecomb-service-center/eventbase/datasource/etcd/tombstone"
-       "servicecomb-service-center/eventbase/datasource/tlsutil"
+       "github.com/apache/servicecomb-service-center/eventbase/datasource"
+       
"github.com/apache/servicecomb-service-center/eventbase/datasource/etcd/task"
+       
"github.com/apache/servicecomb-service-center/eventbase/datasource/etcd/tombstone"
+       
"github.com/apache/servicecomb-service-center/eventbase/datasource/tlsutil"
 )
 
 type Datasource struct {
diff --git a/eventbase/datasource/etcd/etcd_test.go 
b/eventbase/datasource/etcd/etcd_test.go
index 7062c3d..d7c806a 100644
--- a/eventbase/datasource/etcd/etcd_test.go
+++ b/eventbase/datasource/etcd/etcd_test.go
@@ -23,10 +23,12 @@ import (
 
        "github.com/go-chassis/cari/db"
        "github.com/stretchr/testify/assert"
+       // support embedded etcd
+       _ "github.com/little-cui/etcdadpt/embedded"
+       _ "github.com/little-cui/etcdadpt/remote"
 
-       "servicecomb-service-center/eventbase/datasource/etcd"
-       _ "servicecomb-service-center/eventbase/datasource/etcd"
-       "servicecomb-service-center/eventbase/test"
+       "github.com/apache/servicecomb-service-center/eventbase/datasource/etcd"
+       "github.com/apache/servicecomb-service-center/eventbase/test"
 )
 
 func TestNewDatasource(t *testing.T) {
diff --git a/eventbase/datasource/etcd/key/key.go 
b/eventbase/datasource/etcd/key/key.go
index 9464c91..c18750a 100644
--- a/eventbase/datasource/etcd/key/key.go
+++ b/eventbase/datasource/etcd/key/key.go
@@ -43,6 +43,9 @@ func TaskKey(domain, project, taskID string, timestamp int64) 
string {
 }
 
 func TaskList(domain, project string) string {
+       if len(domain) == 0 {
+               return getSyncRootKey()
+       }
        if len(project) == 0 {
                return strings.Join([]string{getSyncRootKey(), domain, ""}, 
split)
        }
@@ -50,6 +53,9 @@ func TaskList(domain, project string) string {
 }
 
 func TombstoneList(domain, project string) string {
+       if len(domain) == 0 {
+               return getTombstoneRootKey()
+       }
        if len(project) == 0 {
                return strings.Join([]string{getTombstoneRootKey(), domain, 
""}, split)
        }
diff --git a/eventbase/datasource/etcd/task/task_dao.go 
b/eventbase/datasource/etcd/task/task_dao.go
index 0e910a2..bc7a68e 100644
--- a/eventbase/datasource/etcd/task/task_dao.go
+++ b/eventbase/datasource/etcd/task/task_dao.go
@@ -25,8 +25,8 @@ import (
        "github.com/go-chassis/openlog"
        "github.com/little-cui/etcdadpt"
 
-       "servicecomb-service-center/eventbase/datasource"
-       "servicecomb-service-center/eventbase/datasource/etcd/key"
+       "github.com/apache/servicecomb-service-center/eventbase/datasource"
+       
"github.com/apache/servicecomb-service-center/eventbase/datasource/etcd/key"
 )
 
 type Dao struct {
@@ -89,13 +89,13 @@ func (d *Dao) Delete(ctx context.Context, tasks 
...*sync.Task) error {
        return nil
 }
 
-func (d Dao) List(ctx context.Context, domain string, project string, options 
...datasource.TaskFindOption) ([]*sync.Task, error) {
+func (d Dao) List(ctx context.Context, options ...datasource.TaskFindOption) 
([]*sync.Task, error) {
        opts := datasource.NewTaskFindOptions()
        for _, o := range options {
                o(&opts)
        }
        tasks := make([]*sync.Task, 0)
-       kvs, _, err := etcdadpt.List(ctx, key.TaskList(domain, project))
+       kvs, _, err := etcdadpt.List(ctx, key.TaskList(opts.Domain, 
opts.Project))
        if err != nil {
                openlog.Error("fail to list task" + err.Error())
                return tasks, err
diff --git a/eventbase/datasource/etcd/task/task_dao_test.go 
b/eventbase/datasource/etcd/task/task_dao_test.go
index 05e83c9..97ac43c 100644
--- a/eventbase/datasource/etcd/task/task_dao_test.go
+++ b/eventbase/datasource/etcd/task/task_dao_test.go
@@ -24,10 +24,13 @@ import (
        "github.com/go-chassis/cari/db"
        "github.com/go-chassis/cari/sync"
        "github.com/stretchr/testify/assert"
+       // support embedded etcd
+       _ "github.com/little-cui/etcdadpt/embedded"
+       _ "github.com/little-cui/etcdadpt/remote"
 
-       "servicecomb-service-center/eventbase/datasource"
-       "servicecomb-service-center/eventbase/datasource/etcd"
-       "servicecomb-service-center/eventbase/test"
+       "github.com/apache/servicecomb-service-center/eventbase/datasource"
+       "github.com/apache/servicecomb-service-center/eventbase/datasource/etcd"
+       "github.com/apache/servicecomb-service-center/eventbase/test"
 )
 
 var ds datasource.DataSource
@@ -110,19 +113,21 @@ func TestTask(t *testing.T) {
        })
 
        t.Run("list task", func(t *testing.T) {
-               t.Run("list task with action ,dataType and status should pass", 
func(t *testing.T) {
+               t.Run("list task with domain, project, action ,dataType and 
status should pass", func(t *testing.T) {
                        opts := []datasource.TaskFindOption{
+                               datasource.WithDomain(task.Domain),
+                               datasource.WithProject(task.Project),
                                datasource.WithAction(task.Action),
                                datasource.WithDataType(task.DataType),
                                datasource.WithStatus(task.Status),
                        }
-                       tasks, err := ds.TaskDao().List(context.Background(), 
task.Domain, task.Project, opts...)
+                       tasks, err := ds.TaskDao().List(context.Background(), 
opts...)
                        assert.NoError(t, err)
                        assert.Equal(t, 1, len(tasks))
                })
 
                t.Run("list task without action ,dataType and status should 
pass", func(t *testing.T) {
-                       tasks, err := ds.TaskDao().List(context.Background(), 
"default", "default")
+                       tasks, err := ds.TaskDao().List(context.Background())
                        assert.NoError(t, err)
                        assert.Equal(t, 3, len(tasks))
                        assert.Equal(t, tasks[0].Timestamp, task.Timestamp)
diff --git a/eventbase/datasource/etcd/tombstone/tombstone_dao.go 
b/eventbase/datasource/etcd/tombstone/tombstone_dao.go
index 12a32bc..26467c5 100644
--- a/eventbase/datasource/etcd/tombstone/tombstone_dao.go
+++ b/eventbase/datasource/etcd/tombstone/tombstone_dao.go
@@ -25,15 +25,15 @@ import (
        "github.com/go-chassis/openlog"
        "github.com/little-cui/etcdadpt"
 
-       "servicecomb-service-center/eventbase/datasource"
-       "servicecomb-service-center/eventbase/datasource/etcd/key"
-       "servicecomb-service-center/eventbase/model"
+       "github.com/apache/servicecomb-service-center/eventbase/datasource"
+       
"github.com/apache/servicecomb-service-center/eventbase/datasource/etcd/key"
+       "github.com/apache/servicecomb-service-center/eventbase/request"
 )
 
 type Dao struct {
 }
 
-func (d *Dao) Get(ctx context.Context, req *model.GetTombstoneRequest) 
(*sync.Tombstone, error) {
+func (d *Dao) Get(ctx context.Context, req *request.GetTombstoneRequest) 
(*sync.Tombstone, error) {
        tombstoneKey := key.TombstoneKey(req.Domain, req.Project, 
req.ResourceType, req.ResourceID)
        kv, err := etcdadpt.Get(ctx, tombstoneKey)
        if err != nil {
@@ -84,13 +84,13 @@ func (d *Dao) Delete(ctx context.Context, tombstones 
...*sync.Tombstone) error {
        return nil
 }
 
-func (d *Dao) List(ctx context.Context, domain string, project string, options 
...datasource.TombstoneFindOption) ([]*sync.Tombstone, error) {
+func (d *Dao) List(ctx context.Context, options 
...datasource.TombstoneFindOption) ([]*sync.Tombstone, error) {
        opts := datasource.NewTombstoneFindOptions()
        for _, o := range options {
                o(&opts)
        }
        tombstones := make([]*sync.Tombstone, 0)
-       kvs, _, err := etcdadpt.List(ctx, key.TombstoneList(domain, project))
+       kvs, _, err := etcdadpt.List(ctx, key.TombstoneList(opts.Domain, 
opts.Project))
        if err != nil {
                openlog.Error("fail to list tombstone" + err.Error())
                return tombstones, err
diff --git a/eventbase/datasource/etcd/tombstone/tombstone_dao_test.go 
b/eventbase/datasource/etcd/tombstone/tombstone_dao_test.go
index 8fe0f91..aedd0f0 100644
--- a/eventbase/datasource/etcd/tombstone/tombstone_dao_test.go
+++ b/eventbase/datasource/etcd/tombstone/tombstone_dao_test.go
@@ -24,11 +24,14 @@ import (
        "github.com/go-chassis/cari/db"
        "github.com/go-chassis/cari/sync"
        "github.com/stretchr/testify/assert"
+       // support embedded etcd
+       _ "github.com/little-cui/etcdadpt/embedded"
+       _ "github.com/little-cui/etcdadpt/remote"
 
-       "servicecomb-service-center/eventbase/datasource"
-       "servicecomb-service-center/eventbase/datasource/etcd"
-       "servicecomb-service-center/eventbase/model"
-       "servicecomb-service-center/eventbase/test"
+       "github.com/apache/servicecomb-service-center/eventbase/datasource"
+       "github.com/apache/servicecomb-service-center/eventbase/datasource/etcd"
+       "github.com/apache/servicecomb-service-center/eventbase/request"
+       "github.com/apache/servicecomb-service-center/eventbase/test"
 )
 
 var ds datasource.DataSource
@@ -72,7 +75,7 @@ func TestTombstone(t *testing.T) {
 
        t.Run("get tombstone", func(t *testing.T) {
                t.Run("get one tombstone should pass", func(t *testing.T) {
-                       req := model.GetTombstoneRequest{
+                       req := request.GetTombstoneRequest{
                                Domain:       tombstoneOne.Domain,
                                Project:      tombstoneOne.Project,
                                ResourceType: tombstoneOne.ResourceType,
@@ -85,12 +88,14 @@ func TestTombstone(t *testing.T) {
        })
 
        t.Run("list tombstone", func(t *testing.T) {
-               t.Run("list tombstone with ResourceType and BeforeTimestamp 
should pass", func(t *testing.T) {
+               t.Run("list tombstone with Domain, Project ,ResourceType and 
BeforeTimestamp should pass", func(t *testing.T) {
                        opts := []datasource.TombstoneFindOption{
+                               datasource.WithTombstoneDomain("default"),
+                               datasource.WithTombstoneDomain("default"),
                                
datasource.WithResourceType(tombstoneOne.ResourceType),
                                datasource.WithBeforeTimestamp(1638171600),
                        }
-                       tombstones, err := 
ds.TombstoneDao().List(context.Background(), "default", "default", opts...)
+                       tombstones, err := 
ds.TombstoneDao().List(context.Background(), opts...)
                        assert.NoError(t, err)
                        assert.Equal(t, 2, len(tombstones))
                        assert.Equal(t, tombstones[0].Timestamp, 
tombstoneOne.Timestamp)
diff --git a/eventbase/datasource/manager.go b/eventbase/datasource/manager.go
index 215d322..06afb9b 100644
--- a/eventbase/datasource/manager.go
+++ b/eventbase/datasource/manager.go
@@ -58,6 +58,7 @@ func Init(c db.Config) error {
                c.Timeout = DefaultTimeout
        }
        dbc := &db.Config{
+               Kind:        c.Kind,
                URI:         c.URI,
                PoolSize:    c.PoolSize,
                SSLEnabled:  c.SSLEnabled,
diff --git a/eventbase/datasource/mongo/client/client.go 
b/eventbase/datasource/mongo/client/client.go
index dea67d7..0998a9c 100644
--- a/eventbase/datasource/mongo/client/client.go
+++ b/eventbase/datasource/mongo/client/client.go
@@ -32,7 +32,7 @@ import (
        "go.mongodb.org/mongo-driver/mongo"
        "go.mongodb.org/mongo-driver/mongo/options"
 
-       dmongo "servicecomb-service-center/eventbase/datasource/mongo"
+       
"github.com/apache/servicecomb-service-center/eventbase/datasource/mongo/model"
 )
 
 const (
@@ -137,7 +137,7 @@ func (mc *MongoClient) newClient(ctx context.Context) (err 
error) {
                }
                return
        }
-       mc.db = mc.client.Database(dmongo.DBName)
+       mc.db = mc.client.Database(model.DBName)
        if mc.db == nil {
                return ErrOpenDbFailed
        }
diff --git a/eventbase/datasource/mongo/types.go 
b/eventbase/datasource/mongo/model/types.go
similarity index 98%
rename from eventbase/datasource/mongo/types.go
rename to eventbase/datasource/mongo/model/types.go
index 96f5c67..04d4421 100644
--- a/eventbase/datasource/mongo/types.go
+++ b/eventbase/datasource/mongo/model/types.go
@@ -15,7 +15,7 @@
  * limitations under the License.
  */
 
-package mongo
+package model
 
 const (
        DBName = "servicecomb"
diff --git a/eventbase/datasource/mongo/mongo.go 
b/eventbase/datasource/mongo/mongo.go
index d156a17..580a6f9 100644
--- a/eventbase/datasource/mongo/mongo.go
+++ b/eventbase/datasource/mongo/mongo.go
@@ -25,10 +25,11 @@ import (
        "go.mongodb.org/mongo-driver/bson"
        "gopkg.in/mgo.v2"
 
-       "servicecomb-service-center/eventbase/datasource"
-       "servicecomb-service-center/eventbase/datasource/mongo/client"
-       "servicecomb-service-center/eventbase/datasource/mongo/task"
-       "servicecomb-service-center/eventbase/datasource/mongo/tombstone"
+       "github.com/apache/servicecomb-service-center/eventbase/datasource"
+       
"github.com/apache/servicecomb-service-center/eventbase/datasource/mongo/client"
+       
"github.com/apache/servicecomb-service-center/eventbase/datasource/mongo/model"
+       
"github.com/apache/servicecomb-service-center/eventbase/datasource/mongo/task"
+       
"github.com/apache/servicecomb-service-center/eventbase/datasource/mongo/tombstone"
 )
 
 type Datasource struct {
@@ -110,32 +111,32 @@ func wrapError(err error, skipMsg ...string) {
 }
 
 func ensureTask(session *mgo.Session) {
-       c := session.DB(DBName).C(CollectionTask)
+       c := session.DB(model.DBName).C(model.CollectionTask)
        err := c.Create(&mgo.CollectionInfo{Validator: bson.M{
-               ColumnTaskID:    bson.M{"$exists": true},
-               ColumnDomain:    bson.M{"$exists": true},
-               ColumnProject:   bson.M{"$exists": true},
-               ColumnTimestamp: bson.M{"$exists": true},
+               model.ColumnTaskID:    bson.M{"$exists": true},
+               model.ColumnDomain:    bson.M{"$exists": true},
+               model.ColumnProject:   bson.M{"$exists": true},
+               model.ColumnTimestamp: bson.M{"$exists": true},
        }})
        wrapError(err)
        err = c.EnsureIndex(mgo.Index{
-               Key:    []string{ColumnDomain, ColumnProject, ColumnTaskID, 
ColumnTimestamp},
+               Key:    []string{model.ColumnDomain, model.ColumnProject, 
model.ColumnTaskID, model.ColumnTimestamp},
                Unique: true,
        })
        wrapError(err)
 }
 
 func ensureTombstone(session *mgo.Session) {
-       c := session.DB(DBName).C(CollectionTombstone)
+       c := session.DB(model.DBName).C(model.CollectionTombstone)
        err := c.Create(&mgo.CollectionInfo{Validator: bson.M{
-               ColumnResourceID:   bson.M{"$exists": true},
-               ColumnDomain:       bson.M{"$exists": true},
-               ColumnProject:      bson.M{"$exists": true},
-               ColumnResourceType: bson.M{"$exists": true},
+               model.ColumnResourceID:   bson.M{"$exists": true},
+               model.ColumnDomain:       bson.M{"$exists": true},
+               model.ColumnProject:      bson.M{"$exists": true},
+               model.ColumnResourceType: bson.M{"$exists": true},
        }})
        wrapError(err)
        err = c.EnsureIndex(mgo.Index{
-               Key:    []string{ColumnDomain, ColumnProject, ColumnResourceID, 
ColumnResourceType},
+               Key:    []string{model.ColumnDomain, model.ColumnProject, 
model.ColumnResourceID, model.ColumnResourceType},
                Unique: true,
        })
        wrapError(err)
diff --git a/eventbase/datasource/mongo/task/task_dao.go 
b/eventbase/datasource/mongo/task/task_dao.go
index 7deef48..c647a51 100644
--- a/eventbase/datasource/mongo/task/task_dao.go
+++ b/eventbase/datasource/mongo/task/task_dao.go
@@ -26,16 +26,16 @@ import (
        "go.mongodb.org/mongo-driver/mongo"
        mopts "go.mongodb.org/mongo-driver/mongo/options"
 
-       "servicecomb-service-center/eventbase/datasource"
-       dmongo "servicecomb-service-center/eventbase/datasource/mongo"
-       "servicecomb-service-center/eventbase/datasource/mongo/client"
+       "github.com/apache/servicecomb-service-center/eventbase/datasource"
+       
"github.com/apache/servicecomb-service-center/eventbase/datasource/mongo/client"
+       
"github.com/apache/servicecomb-service-center/eventbase/datasource/mongo/model"
 )
 
 type Dao struct {
 }
 
 func (d *Dao) Create(ctx context.Context, task *sync.Task) (*sync.Task, error) 
{
-       collection := 
client.GetMongoClient().GetDB().Collection(dmongo.CollectionTask)
+       collection := 
client.GetMongoClient().GetDB().Collection(model.CollectionTask)
        _, err := collection.InsertOne(ctx, task)
        if err != nil {
                openlog.Error("fail to create task" + err.Error())
@@ -45,11 +45,11 @@ func (d *Dao) Create(ctx context.Context, task *sync.Task) 
(*sync.Task, error) {
 }
 
 func (d *Dao) Update(ctx context.Context, task *sync.Task) error {
-       collection := 
client.GetMongoClient().GetDB().Collection(dmongo.CollectionTask)
+       collection := 
client.GetMongoClient().GetDB().Collection(model.CollectionTask)
        result, err := collection.UpdateOne(ctx,
-               bson.M{dmongo.ColumnTaskID: task.TaskID, dmongo.ColumnDomain: 
task.Domain, dmongo.ColumnProject: task.Project, dmongo.ColumnTimestamp: 
task.Timestamp},
+               bson.M{model.ColumnTaskID: task.TaskID, model.ColumnDomain: 
task.Domain, model.ColumnProject: task.Project, model.ColumnTimestamp: 
task.Timestamp},
                bson.D{{Key: "$set", Value: bson.D{
-                       {Key: dmongo.ColumnStatus, Value: task.Status}}},
+                       {Key: model.ColumnStatus, Value: task.Status}}},
                })
        if err != nil {
                openlog.Error("fail to update task" + err.Error())
@@ -68,16 +68,16 @@ func (d *Dao) Delete(ctx context.Context, tasks 
...*sync.Task) error {
        for i, task := range tasks {
                tasksIDs[i] = task.TaskID
                dFilter := bson.D{
-                       {dmongo.ColumnDomain, task.Domain},
-                       {dmongo.ColumnProject, task.Project},
-                       {dmongo.ColumnTaskID, task.TaskID},
-                       {dmongo.ColumnTimestamp, task.Timestamp},
+                       {model.ColumnDomain, task.Domain},
+                       {model.ColumnProject, task.Project},
+                       {model.ColumnTaskID, task.TaskID},
+                       {model.ColumnTimestamp, task.Timestamp},
                }
                filter = append(filter, dFilter)
        }
 
        var deleteFunc = func(sessionContext mongo.SessionContext) error {
-               collection := 
client.GetMongoClient().GetDB().Collection(dmongo.CollectionTask)
+               collection := 
client.GetMongoClient().GetDB().Collection(model.CollectionTask)
                _, err := collection.DeleteMany(sessionContext, bson.M{"$or": 
filter})
                return err
        }
@@ -87,24 +87,30 @@ func (d *Dao) Delete(ctx context.Context, tasks 
...*sync.Task) error {
        }
        return err
 }
-func (d *Dao) List(ctx context.Context, domain string, project string, options 
...datasource.TaskFindOption) ([]*sync.Task, error) {
+func (d *Dao) List(ctx context.Context, options ...datasource.TaskFindOption) 
([]*sync.Task, error) {
        opts := datasource.NewTaskFindOptions()
        for _, o := range options {
                o(&opts)
        }
-       collection := 
client.GetMongoClient().GetDB().Collection(dmongo.CollectionTask)
-       filter := bson.M{dmongo.ColumnDomain: domain, dmongo.ColumnProject: 
project}
+       collection := 
client.GetMongoClient().GetDB().Collection(model.CollectionTask)
+       filter := bson.M{}
+       if opts.Domain != "" {
+               filter[model.ColumnDomain] = opts.Domain
+       }
+       if opts.Project != "" {
+               filter[model.ColumnProject] = opts.Project
+       }
        if opts.Action != "" {
-               filter[dmongo.ColumnAction] = opts.Action
+               filter[model.ColumnAction] = opts.Action
        }
        if opts.DataType != "" {
-               filter[dmongo.ColumnDataType] = opts.DataType
+               filter[model.ColumnDataType] = opts.DataType
        }
        if opts.Status != "" {
-               filter[dmongo.ColumnStatus] = opts.Status
+               filter[model.ColumnStatus] = opts.Status
        }
        opt := mopts.Find().SetSort(map[string]interface{}{
-               dmongo.ColumnTimestamp: 1,
+               model.ColumnTimestamp: 1,
        })
        cur, err := collection.Find(ctx, filter, opt)
        if err != nil {
diff --git a/eventbase/datasource/mongo/task/task_dao_test.go 
b/eventbase/datasource/mongo/task/task_dao_test.go
index 88c4cfb..47c2d7b 100644
--- a/eventbase/datasource/mongo/task/task_dao_test.go
+++ b/eventbase/datasource/mongo/task/task_dao_test.go
@@ -26,9 +26,9 @@ import (
        "github.com/go-chassis/cari/sync"
        "github.com/stretchr/testify/assert"
 
-       "servicecomb-service-center/eventbase/datasource"
-       "servicecomb-service-center/eventbase/datasource/mongo"
-       "servicecomb-service-center/eventbase/test"
+       "github.com/apache/servicecomb-service-center/eventbase/datasource"
+       
"github.com/apache/servicecomb-service-center/eventbase/datasource/mongo"
+       "github.com/apache/servicecomb-service-center/eventbase/test"
 )
 
 var ds datasource.DataSource
@@ -112,19 +112,21 @@ func TestTask(t *testing.T) {
        })
 
        t.Run("list task", func(t *testing.T) {
-               t.Run("list task with action ,dataType and status should pass", 
func(t *testing.T) {
+               t.Run("list task with domain, project, action ,dataType and 
status should pass", func(t *testing.T) {
                        opts := []datasource.TaskFindOption{
+                               datasource.WithDomain(task.Domain),
+                               datasource.WithProject(task.Project),
                                datasource.WithAction(task.Action),
                                datasource.WithDataType(task.DataType),
                                datasource.WithStatus(task.Status),
                        }
-                       tasks, err := ds.TaskDao().List(context.Background(), 
task.Domain, task.Project, opts...)
+                       tasks, err := ds.TaskDao().List(context.Background(), 
opts...)
                        assert.NoError(t, err)
                        assert.Equal(t, 1, len(tasks))
                })
 
                t.Run("list task without action ,dataType and status should 
pass", func(t *testing.T) {
-                       tasks, err := ds.TaskDao().List(context.Background(), 
"default", "default")
+                       tasks, err := ds.TaskDao().List(context.Background())
                        assert.NoError(t, err)
                        assert.Equal(t, 3, len(tasks))
                        assert.Equal(t, tasks[0].Timestamp, task.Timestamp)
diff --git a/eventbase/datasource/mongo/tombstone/tombstone_dao.go 
b/eventbase/datasource/mongo/tombstone/tombstone_dao.go
index 2ef60f2..f83b3eb 100644
--- a/eventbase/datasource/mongo/tombstone/tombstone_dao.go
+++ b/eventbase/datasource/mongo/tombstone/tombstone_dao.go
@@ -25,18 +25,18 @@ import (
        "go.mongodb.org/mongo-driver/bson"
        "go.mongodb.org/mongo-driver/mongo"
 
-       "servicecomb-service-center/eventbase/datasource"
-       dmongo "servicecomb-service-center/eventbase/datasource/mongo"
-       "servicecomb-service-center/eventbase/datasource/mongo/client"
-       "servicecomb-service-center/eventbase/model"
+       "github.com/apache/servicecomb-service-center/eventbase/datasource"
+       
"github.com/apache/servicecomb-service-center/eventbase/datasource/mongo/client"
+       
"github.com/apache/servicecomb-service-center/eventbase/datasource/mongo/model"
+       "github.com/apache/servicecomb-service-center/eventbase/request"
 )
 
 type Dao struct {
 }
 
-func (d *Dao) Get(ctx context.Context, req *model.GetTombstoneRequest) 
(*sync.Tombstone, error) {
-       collection := 
client.GetMongoClient().GetDB().Collection(dmongo.CollectionTombstone)
-       filter := bson.M{dmongo.ColumnDomain: req.Domain, dmongo.ColumnProject: 
req.Project, dmongo.ColumnResourceType: req.ResourceType, 
dmongo.ColumnResourceID: req.ResourceID}
+func (d *Dao) Get(ctx context.Context, req *request.GetTombstoneRequest) 
(*sync.Tombstone, error) {
+       collection := 
client.GetMongoClient().GetDB().Collection(model.CollectionTombstone)
+       filter := bson.M{model.ColumnDomain: req.Domain, model.ColumnProject: 
req.Project, model.ColumnResourceType: req.ResourceType, 
model.ColumnResourceID: req.ResourceID}
        result := collection.FindOne(ctx, filter)
        if result != nil && result.Err() != nil {
                openlog.Error("fail to get tombstone" + result.Err().Error())
@@ -57,7 +57,7 @@ func (d *Dao) Get(ctx context.Context, req 
*model.GetTombstoneRequest) (*sync.To
 }
 
 func (d *Dao) Create(ctx context.Context, tombstone *sync.Tombstone) 
(*sync.Tombstone, error) {
-       collection := 
client.GetMongoClient().GetDB().Collection(dmongo.CollectionTombstone)
+       collection := 
client.GetMongoClient().GetDB().Collection(model.CollectionTombstone)
        _, err := collection.InsertOne(ctx, tombstone)
        if err != nil {
                openlog.Error("fail to create tombstone" + err.Error())
@@ -72,15 +72,15 @@ func (d *Dao) Delete(ctx context.Context, tombstones 
...*sync.Tombstone) error {
        for i, tombstone := range tombstones {
                tombstonesIDs[i] = tombstone.ResourceID
                dFilter := bson.D{
-                       {dmongo.ColumnResourceID, tombstone.ResourceID},
-                       {dmongo.ColumnResourceType, tombstone.ResourceType},
-                       {dmongo.ColumnDomain, tombstone.Domain},
-                       {dmongo.ColumnProject, tombstone.Project},
+                       {model.ColumnResourceID, tombstone.ResourceID},
+                       {model.ColumnResourceType, tombstone.ResourceType},
+                       {model.ColumnDomain, tombstone.Domain},
+                       {model.ColumnProject, tombstone.Project},
                }
                filter = append(filter, dFilter)
        }
        var deleteFunc = func(sessionContext mongo.SessionContext) error {
-               collection := 
client.GetMongoClient().GetDB().Collection(dmongo.CollectionTombstone)
+               collection := 
client.GetMongoClient().GetDB().Collection(model.CollectionTombstone)
                _, err := collection.DeleteMany(sessionContext, bson.M{"$or": 
filter})
                return err
        }
@@ -91,19 +91,24 @@ func (d *Dao) Delete(ctx context.Context, tombstones 
...*sync.Tombstone) error {
        return err
 }
 
-func (d *Dao) List(ctx context.Context, domain string, project string,
-       options ...datasource.TombstoneFindOption) ([]*sync.Tombstone, error) {
+func (d *Dao) List(ctx context.Context, options 
...datasource.TombstoneFindOption) ([]*sync.Tombstone, error) {
        opts := datasource.NewTombstoneFindOptions()
        for _, o := range options {
                o(&opts)
        }
-       collection := 
client.GetMongoClient().GetDB().Collection(dmongo.CollectionTombstone)
-       filter := bson.M{dmongo.ColumnDomain: domain, dmongo.ColumnProject: 
project}
+       collection := 
client.GetMongoClient().GetDB().Collection(model.CollectionTombstone)
+       filter := bson.M{}
+       if opts.Domain != "" {
+               filter[model.ColumnDomain] = opts.Domain
+       }
+       if opts.Project != "" {
+               filter[model.ColumnProject] = opts.Project
+       }
        if opts.ResourceType != "" {
-               filter[dmongo.ColumnResourceType] = opts.ResourceType
+               filter[model.ColumnResourceType] = opts.ResourceType
        }
        if opts.BeforeTimestamp != 0 {
-               filter[dmongo.ColumnTimestamp] = bson.M{"$lte": 
opts.BeforeTimestamp}
+               filter[model.ColumnTimestamp] = bson.M{"$lte": 
opts.BeforeTimestamp}
        }
        cur, err := collection.Find(ctx, filter)
        if err != nil {
diff --git a/eventbase/datasource/mongo/tombstone/tombstone_dao_test.go 
b/eventbase/datasource/mongo/tombstone/tombstone_dao_test.go
index 5f346b1..dccb1da 100644
--- a/eventbase/datasource/mongo/tombstone/tombstone_dao_test.go
+++ b/eventbase/datasource/mongo/tombstone/tombstone_dao_test.go
@@ -26,10 +26,10 @@ import (
        "github.com/go-chassis/cari/sync"
        "github.com/stretchr/testify/assert"
 
-       "servicecomb-service-center/eventbase/datasource"
-       "servicecomb-service-center/eventbase/datasource/mongo"
-       "servicecomb-service-center/eventbase/model"
-       "servicecomb-service-center/eventbase/test"
+       "github.com/apache/servicecomb-service-center/eventbase/datasource"
+       
"github.com/apache/servicecomb-service-center/eventbase/datasource/mongo"
+       "github.com/apache/servicecomb-service-center/eventbase/request"
+       "github.com/apache/servicecomb-service-center/eventbase/test"
 )
 
 var ds datasource.DataSource
@@ -74,7 +74,7 @@ func TestTombstone(t *testing.T) {
 
        t.Run("get tombstone", func(t *testing.T) {
                t.Run("get one tombstone should pass", func(t *testing.T) {
-                       req := model.GetTombstoneRequest{
+                       req := request.GetTombstoneRequest{
                                Domain:       tombstoneOne.Domain,
                                Project:      tombstoneOne.Project,
                                ResourceType: tombstoneOne.ResourceType,
@@ -89,10 +89,12 @@ func TestTombstone(t *testing.T) {
        t.Run("list tombstone", func(t *testing.T) {
                t.Run("list tombstone with ResourceType and BeforeTimestamp 
should pass", func(t *testing.T) {
                        opts := []datasource.TombstoneFindOption{
+                               datasource.WithTombstoneDomain("default"),
+                               datasource.WithTombstoneProject("default"),
                                
datasource.WithResourceType(tombstoneOne.ResourceType),
                                datasource.WithBeforeTimestamp(1638171600),
                        }
-                       tombstones, err := 
ds.TombstoneDao().List(context.Background(), "default", "default", opts...)
+                       tombstones, err := 
ds.TombstoneDao().List(context.Background(), opts...)
                        assert.NoError(t, err)
                        assert.Equal(t, 2, len(tombstones))
                        assert.Equal(t, tombstones[0].Timestamp, 
tombstoneOne.Timestamp)
diff --git a/eventbase/datasource/options.go b/eventbase/datasource/options.go
index 52b2ebc..6c13236 100644
--- a/eventbase/datasource/options.go
+++ b/eventbase/datasource/options.go
@@ -18,12 +18,16 @@
 package datasource
 
 type TaskFindOptions struct {
+       Domain   string
+       Project  string
        Action   string
        Status   string
        DataType string
 }
 
 type TombstoneFindOptions struct {
+       Domain          string
+       Project         string
        ResourceType    string
        BeforeTimestamp int64
 }
@@ -40,6 +44,20 @@ func NewTombstoneFindOptions() TombstoneFindOptions {
        return TombstoneFindOptions{}
 }
 
+// WithDomain find task with domain
+func WithDomain(domain string) TaskFindOption {
+       return func(options *TaskFindOptions) {
+               options.Domain = domain
+       }
+}
+
+// WithProject find task with project
+func WithProject(project string) TaskFindOption {
+       return func(options *TaskFindOptions) {
+               options.Project = project
+       }
+}
+
 // WithAction find task with action
 func WithAction(action string) TaskFindOption {
        return func(options *TaskFindOptions) {
@@ -61,6 +79,20 @@ func WithDataType(dataType string) TaskFindOption {
        }
 }
 
+// WithTombstoneDomain find tombstone with domain
+func WithTombstoneDomain(domain string) TombstoneFindOption {
+       return func(options *TombstoneFindOptions) {
+               options.Domain = domain
+       }
+}
+
+// WithTombstoneProject find tombstone with project
+func WithTombstoneProject(project string) TombstoneFindOption {
+       return func(options *TombstoneFindOptions) {
+               options.Project = project
+       }
+}
+
 // WithResourceType find tombstone with resource type
 func WithResourceType(resourceType string) TombstoneFindOption {
        return func(options *TombstoneFindOptions) {
diff --git a/eventbase/datasource/task.go b/eventbase/datasource/task.go
index f299109..05b2f06 100644
--- a/eventbase/datasource/task.go
+++ b/eventbase/datasource/task.go
@@ -29,5 +29,5 @@ type TaskDao interface {
        Create(ctx context.Context, task *sync.Task) (*sync.Task, error)
        Update(ctx context.Context, task *sync.Task) error
        Delete(ctx context.Context, tasks ...*sync.Task) error
-       List(ctx context.Context, domain string, project string, options 
...TaskFindOption) ([]*sync.Task, error)
+       List(ctx context.Context, options ...TaskFindOption) ([]*sync.Task, 
error)
 }
diff --git a/eventbase/datasource/tlsutil/tlsutil_test.go 
b/eventbase/datasource/tlsutil/tlsutil_test.go
index f3daf26..fd99c05 100644
--- a/eventbase/datasource/tlsutil/tlsutil_test.go
+++ b/eventbase/datasource/tlsutil/tlsutil_test.go
@@ -26,7 +26,7 @@ import (
        _ "github.com/go-chassis/go-chassis/v2/security/cipher/plugins/plain"
        "github.com/stretchr/testify/assert"
 
-       "servicecomb-service-center/eventbase/datasource/tlsutil"
+       
"github.com/apache/servicecomb-service-center/eventbase/datasource/tlsutil"
 )
 
 const sslRoot = "./../../../examples/service_center/ssl/"
diff --git a/eventbase/datasource/tombstone.go 
b/eventbase/datasource/tombstone.go
index 97433e3..ee1e924 100644
--- a/eventbase/datasource/tombstone.go
+++ b/eventbase/datasource/tombstone.go
@@ -22,14 +22,14 @@ import (
 
        "github.com/go-chassis/cari/sync"
 
-       "servicecomb-service-center/eventbase/model"
+       "github.com/apache/servicecomb-service-center/eventbase/request"
 )
 
 // TombstoneDao provide api of Tombstone entity
 type TombstoneDao interface {
-       Get(ctx context.Context, req *model.GetTombstoneRequest) 
(*sync.Tombstone, error)
+       Get(ctx context.Context, req *request.GetTombstoneRequest) 
(*sync.Tombstone, error)
        // Create func is used for ut
        Create(ctx context.Context, tombstone *sync.Tombstone) 
(*sync.Tombstone, error)
        Delete(ctx context.Context, tombstones ...*sync.Tombstone) error
-       List(ctx context.Context, domain string, project string, options 
...TombstoneFindOption) ([]*sync.Tombstone, error)
+       List(ctx context.Context, options ...TombstoneFindOption) 
([]*sync.Tombstone, error)
 }
diff --git a/eventbase/go.mod b/eventbase/go.mod
index b8492b9..4a659fe 100644
--- a/eventbase/go.mod
+++ b/eventbase/go.mod
@@ -1,4 +1,4 @@
-module servicecomb-service-center/eventbase
+module github.com/apache/servicecomb-service-center/eventbase
 
 require (
        github.com/go-chassis/cari v0.5.1-0.20211208092532-78a52aa9d52e
diff --git a/eventbase/model/tombstone_request.go 
b/eventbase/request/tombstone_request.go
similarity index 56%
rename from eventbase/model/tombstone_request.go
rename to eventbase/request/tombstone_request.go
index 1e77f62..63d1cb0 100644
--- a/eventbase/model/tombstone_request.go
+++ b/eventbase/request/tombstone_request.go
@@ -15,7 +15,7 @@
  * limitations under the License.
  */
 
-package model
+package request
 
 // GetTombstoneRequest contains tombstone get request params
 type GetTombstoneRequest struct {
@@ -24,3 +24,20 @@ type GetTombstoneRequest struct {
        ResourceType string `json:"resource_type,omitempty" 
yaml:"resource_type,omitempty"`
        ResourceID   string `json:"resource_id,omitempty" 
yaml:"resource_id,omitempty"`
 }
+
+// ListTaskRequest contains task list request params
+type ListTaskRequest struct {
+       Domain       string `json:"domain,omitempty" yaml:"domain,omitempty"`
+       Project      string `json:"project,omitempty" yaml:"project,omitempty"`
+       TaskAction   string `json:"task_action,omitempty" 
yaml:"task_action,omitempty"`
+       TaskStatus   string `json:"task_status,omitempty" 
yaml:"task_status,omitempty"`
+       TaskDataType string `json:"task_data_type,omitempty" 
yaml:"task_data_type,omitempty"`
+}
+
+// ListTombstoneRequest contains tombstone list request params
+type ListTombstoneRequest struct {
+       Domain          string `json:"domain,omitempty" yaml:"domain,omitempty"`
+       Project         string `json:"project,omitempty" 
yaml:"project,omitempty"`
+       ResourceType    string `json:"resource_type,omitempty" 
yaml:"resource_type,omitempty"`
+       BeforeTimestamp int64  `json:"before_timestamp,omitempty" 
yaml:"before_timestamp,omitempty"`
+}
\ No newline at end of file
diff --git a/eventbase/datasource/tombstone.go 
b/eventbase/service/task/task_svc.go
similarity index 50%
copy from eventbase/datasource/tombstone.go
copy to eventbase/service/task/task_svc.go
index 97433e3..d35bfbf 100644
--- a/eventbase/datasource/tombstone.go
+++ b/eventbase/service/task/task_svc.go
@@ -15,21 +15,36 @@
  * limitations under the License.
  */
 
-package datasource
+package task
 
 import (
        "context"
 
        "github.com/go-chassis/cari/sync"
 
-       "servicecomb-service-center/eventbase/model"
+       "github.com/apache/servicecomb-service-center/eventbase/datasource"
+       "github.com/apache/servicecomb-service-center/eventbase/request"
 )
 
-// TombstoneDao provide api of Tombstone entity
-type TombstoneDao interface {
-       Get(ctx context.Context, req *model.GetTombstoneRequest) 
(*sync.Tombstone, error)
-       // Create func is used for ut
-       Create(ctx context.Context, tombstone *sync.Tombstone) 
(*sync.Tombstone, error)
-       Delete(ctx context.Context, tombstones ...*sync.Tombstone) error
-       List(ctx context.Context, domain string, project string, options 
...TombstoneFindOption) ([]*sync.Tombstone, error)
+func Delete(ctx context.Context, tasks ...*sync.Task) error {
+       return datasource.GetTaskDao().Delete(ctx, tasks...)
+}
+
+func Update(ctx context.Context, task *sync.Task) error {
+       return datasource.GetTaskDao().Update(ctx, task)
+}
+
+func List(ctx context.Context, request *request.ListTaskRequest) 
([]*sync.Task, error) {
+       opts := []datasource.TaskFindOption{
+               datasource.WithDomain(request.Domain),
+               datasource.WithProject(request.Project),
+               datasource.WithAction(request.TaskAction),
+               datasource.WithDataType(request.TaskDataType),
+               datasource.WithStatus(request.TaskStatus),
+       }
+       tasks, err := datasource.GetTaskDao().List(ctx, opts...)
+       if err != nil {
+               return nil, err
+       }
+       return tasks, nil
 }
diff --git a/eventbase/datasource/tombstone.go 
b/eventbase/service/tombstone/tombstone_svc.go
similarity index 50%
copy from eventbase/datasource/tombstone.go
copy to eventbase/service/tombstone/tombstone_svc.go
index 97433e3..17e7bbb 100644
--- a/eventbase/datasource/tombstone.go
+++ b/eventbase/service/tombstone/tombstone_svc.go
@@ -15,21 +15,31 @@
  * limitations under the License.
  */
 
-package datasource
+package tombstone
 
 import (
        "context"
 
        "github.com/go-chassis/cari/sync"
 
-       "servicecomb-service-center/eventbase/model"
+       "github.com/apache/servicecomb-service-center/eventbase/datasource"
+       "github.com/apache/servicecomb-service-center/eventbase/request"
 )
 
-// TombstoneDao provide api of Tombstone entity
-type TombstoneDao interface {
-       Get(ctx context.Context, req *model.GetTombstoneRequest) 
(*sync.Tombstone, error)
-       // Create func is used for ut
-       Create(ctx context.Context, tombstone *sync.Tombstone) 
(*sync.Tombstone, error)
-       Delete(ctx context.Context, tombstones ...*sync.Tombstone) error
-       List(ctx context.Context, domain string, project string, options 
...TombstoneFindOption) ([]*sync.Tombstone, error)
+func Get(ctx context.Context, req *request.GetTombstoneRequest) 
(*sync.Tombstone, error) {
+       return datasource.GetTombstoneDao().Get(ctx, req)
+}
+
+func Delete(ctx context.Context, tombstones ...*sync.Tombstone) error {
+       return datasource.GetTombstoneDao().Delete(ctx, tombstones...)
+}
+
+func List(ctx context.Context, request *request.ListTombstoneRequest) 
([]*sync.Tombstone, error) {
+       opts := []datasource.TombstoneFindOption{
+               datasource.WithTombstoneDomain(request.Domain),
+               datasource.WithTombstoneProject(request.Project),
+               datasource.WithResourceType(request.ResourceType),
+               datasource.WithBeforeTimestamp(request.BeforeTimestamp),
+       }
+       return datasource.GetTombstoneDao().List(ctx, opts...)
 }

Reply via email to