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...)
}