mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-24 16:03:43 +08:00
Merge pull request #457 from yousong/bugfix/yousong-enhance-db-tasks
Bugfix/yousong enhance db tasks
This commit is contained in:
Generated
+2
-2
@@ -1762,11 +1762,11 @@
|
||||
|
||||
[[projects]]
|
||||
branch = "master"
|
||||
digest = "1:ea481ba82e96a2a5c70d259fc0766ea18d7b6f59c74d3cac38a456e8c6c3e828"
|
||||
digest = "1:29f0a99ba079309ef8c82def5b9c452d0871837780ee702d7ae7c1ebe47478b0"
|
||||
name = "yunion.io/x/sqlchemy"
|
||||
packages = ["."]
|
||||
pruneopts = "UT"
|
||||
revision = "bf5e3b0d446dae4e9bf66cc82bfbb8c7d6b56dfd"
|
||||
revision = "5e43c9cfbb8eb2870882b87d3ab429062d69be2b"
|
||||
|
||||
[[projects]]
|
||||
branch = "master"
|
||||
|
||||
@@ -65,7 +65,6 @@ var TaskManager *STaskManager
|
||||
|
||||
func init() {
|
||||
TaskManager = &STaskManager{SResourceBaseManager: db.NewResourceBaseManager(STask{}, "tasks_tbl", "task", "tasks")}
|
||||
TaskManager.TableSpec().AddIndex(true, "created_at", "stage", "obj_id", "obj_name")
|
||||
}
|
||||
|
||||
type STask struct {
|
||||
@@ -73,10 +72,10 @@ type STask struct {
|
||||
|
||||
Id string `width:"36" charset:"ascii" primary:"true" list:"user"` // Column(VARCHAR(36, charset='ascii'), primary_key=True, default=get_uuid)
|
||||
|
||||
ObjName string `width:"128" charset:"utf8" nullable:"false" list:"user"` // Column(VARCHAR(128, charset='utf8'), nullable=False)
|
||||
ObjId string `width:"128" charset:"ascii" nullable:"false" list:"user"` // Column(VARCHAR(ID_LENGTH, charset='ascii'), nullable=False)
|
||||
TaskName string `width:"64" charset:"ascii" nullable:"false" list:"user"` // Column(VARCHAR(64, charset='ascii'), nullable=False)
|
||||
UserCred mcclient.TokenCredential `width:"1024" charset:"ascii" nullable:"false" get:"user"` // Column(VARCHAR(1024, charset='ascii'), nullable=False)
|
||||
ObjName string `width:"128" charset:"utf8" nullable:"false" list:"user"` // Column(VARCHAR(128, charset='utf8'), nullable=False)
|
||||
ObjId string `width:"128" charset:"ascii" nullable:"false" list:"user" index:"true"` // Column(VARCHAR(ID_LENGTH, charset='ascii'), nullable=False)
|
||||
TaskName string `width:"64" charset:"ascii" nullable:"false" list:"user"` // Column(VARCHAR(64, charset='ascii'), nullable=False)
|
||||
UserCred mcclient.TokenCredential `width:"1024" charset:"ascii" nullable:"false" get:"user"` // Column(VARCHAR(1024, charset='ascii'), nullable=False)
|
||||
// OwnerCred string `width:"512" charset:"ascii" nullable:"true"` // Column(VARCHAR(512, charset='ascii'), nullable=True)
|
||||
Params *jsonutils.JSONDict `charset:"ascii" length:"medium" nullable:"false" get:"user"` // Column(MEDIUMTEXT(charset='ascii'), nullable=False)
|
||||
|
||||
@@ -706,48 +705,53 @@ func (task *STask) GetStartTime() time.Time {
|
||||
}
|
||||
|
||||
func (manager *STaskManager) QueryTasksOfObject(obj db.IStandaloneModel, since time.Time, isOpen *bool) *sqlchemy.SQuery {
|
||||
subq1 := manager.Query("id")
|
||||
subq1 = subq1.Equals("obj_name", obj.Keyword())
|
||||
subq1 = subq1.Equals("obj_id", obj.GetId())
|
||||
if !since.IsZero() {
|
||||
subq1 = subq1.GE("created_at", since)
|
||||
}
|
||||
if isOpen != nil {
|
||||
if *isOpen {
|
||||
subq1 = subq1.Filter(sqlchemy.NOT(
|
||||
sqlchemy.In(subq1.Field("stage"), []string{"complete", "failed"}),
|
||||
))
|
||||
} else if !*isOpen {
|
||||
subq1 = subq1.In("stage", []string{"complete", "failed"})
|
||||
subq1 := manager.Query()
|
||||
{
|
||||
subq1 = subq1.Equals("obj_id", obj.GetId())
|
||||
subq1 = subq1.Equals("obj_name", obj.Keyword())
|
||||
if !since.IsZero() {
|
||||
subq1 = subq1.GE("created_at", since)
|
||||
}
|
||||
if isOpen != nil {
|
||||
if *isOpen {
|
||||
subq1 = subq1.Filter(sqlchemy.NOT(
|
||||
sqlchemy.In(subq1.Field("stage"), []string{"complete", "failed"}),
|
||||
))
|
||||
} else if !*isOpen {
|
||||
subq1 = subq1.In("stage", []string{"complete", "failed"})
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
taskObjs := TaskObjectManager.Query().SubQuery()
|
||||
subq2 := manager.Query("id").Distinct()
|
||||
subq2 = subq2.Join(taskObjs, sqlchemy.Equals(taskObjs.Field("task_id"), subq2.Field("id")))
|
||||
subq2 = subq2.Filter(sqlchemy.Equals(subq2.Field("obj_id"), MULTI_OBJECTS_ID))
|
||||
subq2 = subq2.Filter(sqlchemy.Equals(subq2.Field("obj_name"), obj.Keyword()))
|
||||
subq2 = subq2.Filter(sqlchemy.Equals(taskObjs.Field("obj_id"), obj.GetId()))
|
||||
if !since.IsZero() {
|
||||
subq2 = subq2.Filter(sqlchemy.GE(subq2.Field("created_at"), since))
|
||||
}
|
||||
if isOpen != nil {
|
||||
if *isOpen {
|
||||
subq2 = subq2.Filter(sqlchemy.NOT(
|
||||
sqlchemy.In(subq2.Field("stage"), []string{"complete", "failed"}),
|
||||
))
|
||||
} else if !*isOpen {
|
||||
subq2 = subq2.In("stage", []string{"complete", "failed"})
|
||||
subq2 := manager.Query()
|
||||
{
|
||||
taskObjs := TaskObjectManager.TableSpec().Instance()
|
||||
subq2 = subq2.Join(taskObjs, sqlchemy.AND(
|
||||
sqlchemy.Equals(taskObjs.Field("task_id"), subq2.Field("id")),
|
||||
sqlchemy.Equals(taskObjs.Field("obj_id"), obj.GetId()),
|
||||
))
|
||||
subq2 = subq2.Filter(sqlchemy.Equals(subq2.Field("obj_id"), MULTI_OBJECTS_ID))
|
||||
subq2 = subq2.Filter(sqlchemy.Equals(subq2.Field("obj_name"), obj.Keyword()))
|
||||
if !since.IsZero() {
|
||||
subq2 = subq2.Filter(sqlchemy.GE(subq2.Field("created_at"), since))
|
||||
}
|
||||
if isOpen != nil {
|
||||
if *isOpen {
|
||||
subq2 = subq2.Filter(sqlchemy.NOT(
|
||||
sqlchemy.In(subq2.Field("stage"), []string{"complete", "failed"}),
|
||||
))
|
||||
} else if !*isOpen {
|
||||
subq2 = subq2.In("stage", []string{"complete", "failed"})
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
q := manager.Query()
|
||||
q = q.Filter(sqlchemy.OR(
|
||||
sqlchemy.In(q.Field("id"), subq1.SubQuery()),
|
||||
sqlchemy.In(q.Field("id"), subq2.SubQuery()),
|
||||
))
|
||||
q = q.Desc("created_at")
|
||||
// subq1 and subq2 do not intersect for the fact that they have
|
||||
// different condition on tasks_tbl.obj_id field
|
||||
uq := sqlchemy.Union(subq1, subq2)
|
||||
uq = uq.Desc("created_at")
|
||||
|
||||
q := uq.SubQuery().Query()
|
||||
return q
|
||||
}
|
||||
|
||||
|
||||
-2
@@ -2,7 +2,6 @@ package sqlchemy
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"sort"
|
||||
"strings"
|
||||
)
|
||||
|
||||
@@ -51,7 +50,6 @@ func (index *STableIndex) QuotedColumns() []string {
|
||||
}
|
||||
|
||||
func (ts *STableSpec) AddIndex(unique bool, cols ...string) bool {
|
||||
sort.Sort(TColumnNames(cols))
|
||||
for i := 0; i < len(ts.indexes); i += 1 {
|
||||
if ts.indexes[i].IsIdentical(cols...) {
|
||||
return false
|
||||
|
||||
+12
-16
@@ -272,15 +272,13 @@ func queryString(tq *SQuery) string {
|
||||
}
|
||||
buf.WriteString(" FROM ")
|
||||
buf.WriteString(fmt.Sprintf("%s AS `%s`", tq.from.Expression(), tq.from.Alias()))
|
||||
if tq.joins != nil && len(tq.joins) > 0 {
|
||||
for _, join := range tq.joins {
|
||||
buf.WriteByte(' ')
|
||||
buf.WriteString(string(join.jointype))
|
||||
buf.WriteByte(' ')
|
||||
buf.WriteString(fmt.Sprintf("%s AS `%s`", join.from.Expression(), join.from.Alias()))
|
||||
buf.WriteString(" ON ")
|
||||
buf.WriteString(join.condition.WhereClause())
|
||||
}
|
||||
for _, join := range tq.joins {
|
||||
buf.WriteByte(' ')
|
||||
buf.WriteString(string(join.jointype))
|
||||
buf.WriteByte(' ')
|
||||
buf.WriteString(fmt.Sprintf("%s AS `%s`", join.from.Expression(), join.from.Alias()))
|
||||
buf.WriteString(" ON ")
|
||||
buf.WriteString(join.condition.WhereClause())
|
||||
}
|
||||
if tq.where != nil {
|
||||
buf.WriteString(" WHERE ")
|
||||
@@ -345,13 +343,11 @@ func (tq *SQuery) Variables() []interface{} {
|
||||
fromvars = tq.from.Variables()
|
||||
vars = append(vars, fromvars...)
|
||||
}
|
||||
if tq.joins != nil && len(tq.joins) > 0 {
|
||||
for _, join := range tq.joins {
|
||||
fromvars = join.from.Variables()
|
||||
vars = append(vars, fromvars...)
|
||||
fromvars = join.condition.Variables()
|
||||
vars = append(vars, fromvars...)
|
||||
}
|
||||
for _, join := range tq.joins {
|
||||
fromvars = join.from.Variables()
|
||||
vars = append(vars, fromvars...)
|
||||
fromvars = join.condition.Variables()
|
||||
vars = append(vars, fromvars...)
|
||||
}
|
||||
if tq.where != nil {
|
||||
fromvars = tq.where.Variables()
|
||||
|
||||
+1
@@ -232,6 +232,7 @@ func (ts *STableSpec) SyncSQL() []string {
|
||||
|
||||
for _, idx := range removeIndexes {
|
||||
sql := fmt.Sprintf("DROP INDEX `%s` ON `%s`", idx.name, ts.name)
|
||||
ret = append(ret, sql)
|
||||
log.Infof(sql)
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user