mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-24 16:03:43 +08:00
feature: archived task show (#22309)
Co-authored-by: Qiu Jian <qiujian@yunionyun.com>
This commit is contained in:
@@ -26,6 +26,7 @@ import (
|
||||
"yunion.io/x/sqlchemy"
|
||||
"yunion.io/x/sqlchemy/backends/clickhouse"
|
||||
|
||||
"yunion.io/x/onecloud/pkg/apis"
|
||||
"yunion.io/x/onecloud/pkg/cloudcommon/consts"
|
||||
"yunion.io/x/onecloud/pkg/httperrors"
|
||||
"yunion.io/x/onecloud/pkg/mcclient"
|
||||
@@ -145,3 +146,18 @@ func (l *SLogBase) ValidateUpdateData(
|
||||
) (jsonutils.JSONObject, error) {
|
||||
return nil, errors.Wrap(httperrors.ErrForbidden, "not allow")
|
||||
}
|
||||
|
||||
// 操作日志列表
|
||||
func (manager *SLogBaseManager) ListItemFilter(
|
||||
ctx context.Context,
|
||||
q *sqlchemy.SQuery,
|
||||
userCred mcclient.TokenCredential,
|
||||
input apis.LogBaseListInput,
|
||||
) (*sqlchemy.SQuery, error) {
|
||||
var err error
|
||||
q, err = manager.SModelBaseManager.ListItemFilter(ctx, q, userCred, input.ModelBaseListInput)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "SModelBaseManager.ListItemFilter")
|
||||
}
|
||||
return q, nil
|
||||
}
|
||||
|
||||
@@ -23,4 +23,10 @@ import (
|
||||
func AddTaskHandler(prefix string, app *appsrv.Application) {
|
||||
handler := db.NewModelHandler(TaskManager)
|
||||
dispatcher.AddModelDispatcher(prefix, app, handler)
|
||||
|
||||
{
|
||||
initArchivedTaskManager()
|
||||
archiveHandler := db.NewModelHandler(ArchivedTaskManager)
|
||||
dispatcher.AddModelDispatcher(prefix, app, archiveHandler)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -22,9 +22,12 @@ import (
|
||||
"yunion.io/x/pkg/util/rbacscope"
|
||||
"yunion.io/x/sqlchemy"
|
||||
|
||||
"yunion.io/x/onecloud/pkg/apis"
|
||||
"yunion.io/x/onecloud/pkg/cloudcommon/consts"
|
||||
"yunion.io/x/onecloud/pkg/cloudcommon/db"
|
||||
"yunion.io/x/onecloud/pkg/httperrors"
|
||||
"yunion.io/x/onecloud/pkg/mcclient"
|
||||
"yunion.io/x/onecloud/pkg/util/stringutils2"
|
||||
)
|
||||
|
||||
type SArchivedTaskManager struct {
|
||||
@@ -34,13 +37,13 @@ type SArchivedTaskManager struct {
|
||||
|
||||
var ArchivedTaskManager *SArchivedTaskManager
|
||||
|
||||
func InitArchivedTaskManager() {
|
||||
func initArchivedTaskManager() {
|
||||
ArchivedTaskManager = &SArchivedTaskManager{
|
||||
SLogBaseManager: db.NewLogBaseManager(
|
||||
SArchivedTask{},
|
||||
"archived_tasks_tbl",
|
||||
"achivedtask",
|
||||
"achivedtasks",
|
||||
"archivedtask",
|
||||
"archivedtasks",
|
||||
"start_at",
|
||||
consts.OpsLogWithClickhouse,
|
||||
),
|
||||
@@ -100,3 +103,140 @@ func (manager *SArchivedTaskManager) FilterByOwner(ctx context.Context, q *sqlch
|
||||
func (manager *SArchivedTaskManager) FetchOwnerId(ctx context.Context, data jsonutils.JSONObject) (mcclient.IIdentityProvider, error) {
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
// 操作日志列表
|
||||
func (manager *SArchivedTaskManager) ListItemFilter(
|
||||
ctx context.Context,
|
||||
q *sqlchemy.SQuery,
|
||||
userCred mcclient.TokenCredential,
|
||||
input apis.ArchivedTaskListInput,
|
||||
) (*sqlchemy.SQuery, error) {
|
||||
var err error
|
||||
|
||||
q, err = manager.SLogBaseManager.ListItemFilter(ctx, q, userCred, input.LogBaseListInput)
|
||||
if err != nil {
|
||||
return q, errors.Wrap(err, "SLogBaseManager.ListItemFilter")
|
||||
}
|
||||
q, err = manager.SStatusResourceBaseManager.ListItemFilter(ctx, q, userCred, input.StatusResourceBaseListInput)
|
||||
if err != nil {
|
||||
return q, errors.Wrap(err, "SStatusResourceBaseManager.ListItemFilter")
|
||||
}
|
||||
|
||||
if len(input.TaskId) > 0 {
|
||||
q = q.In("task_id", input.TaskId)
|
||||
}
|
||||
|
||||
if len(input.ObjId) > 0 {
|
||||
if len(input.ObjId) == 1 {
|
||||
q = q.Contains("obj_ids", input.ObjId[0])
|
||||
} else {
|
||||
filters := make([]sqlchemy.ICondition, 0)
|
||||
for i := range input.ObjId {
|
||||
filters = append(filters, sqlchemy.Contains(q.Field("obj_ids"), input.ObjId[i]))
|
||||
}
|
||||
q = q.Filter(sqlchemy.OR(filters...))
|
||||
}
|
||||
}
|
||||
|
||||
if len(input.ObjName) > 0 {
|
||||
if len(input.ObjName) == 1 {
|
||||
q = q.Contains("obj_names", input.ObjName[0])
|
||||
} else {
|
||||
filters := make([]sqlchemy.ICondition, 0)
|
||||
for i := range input.ObjName {
|
||||
filters = append(filters, sqlchemy.Contains(q.Field("obj_names"), input.ObjName[i]))
|
||||
}
|
||||
q = q.Filter(sqlchemy.OR(filters...))
|
||||
}
|
||||
}
|
||||
|
||||
if len(input.ObjType) > 0 {
|
||||
q = q.In("obj_type", input.ObjType)
|
||||
}
|
||||
|
||||
if len(input.TaskName) > 0 {
|
||||
q = q.In("task_name", input.TaskName)
|
||||
}
|
||||
|
||||
if len(input.Stage) > 0 {
|
||||
q = q.In("stage", input.Stage)
|
||||
}
|
||||
|
||||
if len(input.NotStage) > 0 {
|
||||
q = q.NotIn("stage", input.NotStage)
|
||||
}
|
||||
|
||||
if len(input.ParentId) > 0 {
|
||||
q = q.In("parent_task_id", input.ParentId)
|
||||
}
|
||||
|
||||
if input.IsRoot != nil {
|
||||
if *input.IsRoot {
|
||||
q = q.IsNullOrEmpty("parent_task_id")
|
||||
} else {
|
||||
q = q.IsNotEmpty("parent_task_id")
|
||||
}
|
||||
}
|
||||
|
||||
if len(input.ParentTaskId) > 0 {
|
||||
q = q.Equals("parent_task_id", input.ParentTaskId)
|
||||
}
|
||||
|
||||
return q, nil
|
||||
}
|
||||
|
||||
func (manager *SArchivedTaskManager) ListItemExportKeys(ctx context.Context, q *sqlchemy.SQuery, userCred mcclient.TokenCredential, keys stringutils2.SSortedStrings) (*sqlchemy.SQuery, error) {
|
||||
var err error
|
||||
q, err = manager.SModelBaseManager.ListItemExportKeys(ctx, q, userCred, keys)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "SModelBaseManager.ListItemExportKeys")
|
||||
}
|
||||
// q, err = manager.SProjectizedResourceBaseManager.ListItemExportKeys(ctx, q, userCred, keys)
|
||||
// if err != nil {
|
||||
// return nil, errors.Wrap(err, "SProjectizedResourceBaseManager.ListItemExportKeys")
|
||||
// }
|
||||
return q, nil
|
||||
}
|
||||
|
||||
func (manager *SArchivedTaskManager) QueryDistinctExtraField(q *sqlchemy.SQuery, field string) (*sqlchemy.SQuery, error) {
|
||||
var err error
|
||||
q, err = manager.SModelBaseManager.QueryDistinctExtraField(q, field)
|
||||
if err == nil {
|
||||
return q, nil
|
||||
}
|
||||
// q, err = manager.SProjectizedResourceBaseManager.QueryDistinctExtraField(q, field)
|
||||
// if err == nil {
|
||||
// return q, nil
|
||||
// }
|
||||
return q, httperrors.ErrNotFound
|
||||
}
|
||||
|
||||
func (manager *SArchivedTaskManager) ResourceScope() rbacscope.TRbacScope {
|
||||
return rbacscope.ScopeProject
|
||||
}
|
||||
|
||||
/*func (manager *SArchivedTaskManager) FetchCustomizeColumns(
|
||||
ctx context.Context,
|
||||
userCred mcclient.TokenCredential,
|
||||
query jsonutils.JSONObject,
|
||||
objs []interface{},
|
||||
fields stringutils2.SSortedStrings,
|
||||
isList bool,
|
||||
) []apis.TaskDetails {
|
||||
rows := make([]apis.TaskDetails, len(objs))
|
||||
return rows
|
||||
}*/
|
||||
|
||||
func (manager *SArchivedTaskManager) OrderByExtraFields(
|
||||
ctx context.Context,
|
||||
q *sqlchemy.SQuery,
|
||||
userCred mcclient.TokenCredential,
|
||||
query apis.ArchivedTaskListInput,
|
||||
) (*sqlchemy.SQuery, error) {
|
||||
var err error
|
||||
q, err = manager.SModelBaseManager.OrderByExtraFields(ctx, q, userCred, query.ModelBaseListInput)
|
||||
if err != nil {
|
||||
return q, errors.Wrap(err, "SModelBaseManager.OrderByExtraField")
|
||||
}
|
||||
return q, nil
|
||||
}
|
||||
|
||||
@@ -1115,6 +1115,10 @@ func (manager *STaskManager) ListItemFilter(
|
||||
}
|
||||
}
|
||||
|
||||
if len(input.ParentTaskId) > 0 {
|
||||
q = q.Equals("parent_task_id", input.ParentTaskId)
|
||||
}
|
||||
|
||||
if input.SubTask != nil && *input.SubTask {
|
||||
subSQFunc := func(status string, cntField string) *sqlchemy.SSubQuery {
|
||||
subQ := SubTaskManager.Query()
|
||||
@@ -1275,9 +1279,7 @@ func (manager *STaskManager) InitializeData() error {
|
||||
}
|
||||
|
||||
func (manager *STaskManager) failTimeoutTasks() error {
|
||||
// failed unfinished tasks 24 hours ago
|
||||
q := manager.Query().NotIn("stage", []string{TASK_STAGE_FAILED, TASK_STAGE_COMPLETE})
|
||||
// q = q.LT("created_at", time.Now().Add(-24*time.Hour))
|
||||
|
||||
tasks := make([]STask, 0)
|
||||
err := db.FetchModelObjects(manager, q, &tasks)
|
||||
@@ -1470,3 +1472,49 @@ func (manager *STaskManager) TaskCleanupJob(ctx context.Context, userCred mcclie
|
||||
}
|
||||
log.Infof("TaskCleanupJob migrate %d tasks, takes %f seconds, batch limit=%d threshold hours=%d", count, time.Since(taskStart).Seconds(), consts.TaskArchiveBatchLimit(), consts.TaskArchiveThresholdHours())
|
||||
}
|
||||
|
||||
func (task *STask) PerformCancel(
|
||||
ctx context.Context,
|
||||
userCred mcclient.TokenCredential,
|
||||
query jsonutils.JSONObject,
|
||||
input apis.TaskCancelInput,
|
||||
) (jsonutils.JSONObject, error) {
|
||||
if utils.IsInArray(task.Stage, []string{TASK_STAGE_FAILED, TASK_STAGE_COMPLETE}) {
|
||||
return nil, errors.Wrapf(errors.ErrInvalidStatus, "cannot cancel stage in %s", task.Stage)
|
||||
}
|
||||
err := task.cancel(ctx)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "cancel")
|
||||
}
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
func (task *STask) fetchSubTasks() ([]STask, error) {
|
||||
q := task.GetModelManager().Query().Equals("parent_task_id", task.Id)
|
||||
|
||||
tasks := make([]STask, 0)
|
||||
err := db.FetchModelObjects(task.GetModelManager(), q, &tasks)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "db.FetchModelObjects")
|
||||
}
|
||||
return tasks, nil
|
||||
}
|
||||
|
||||
func (task *STask) cancel(ctx context.Context) error {
|
||||
if utils.IsInArray(task.Stage, []string{TASK_STAGE_FAILED, TASK_STAGE_COMPLETE}) {
|
||||
return nil
|
||||
}
|
||||
subtasks, err := task.fetchSubTasks()
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "fetchSubTasks")
|
||||
}
|
||||
for i := range subtasks {
|
||||
err := subtasks[i].cancel(ctx)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "cancelTask")
|
||||
}
|
||||
}
|
||||
reason := jsonutils.NewString("cancel")
|
||||
task.SetStageFailed(ctx, reason)
|
||||
return nil
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user