Skip to content

Commit

Permalink
增加可配置的定时任务功能
Browse files Browse the repository at this point in the history
  • Loading branch information
dushixiang committed Feb 28, 2021
1 parent 5ea00d8 commit 1b8ecef
Show file tree
Hide file tree
Showing 10 changed files with 862 additions and 83 deletions.
6 changes: 5 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
# Next Terminal

你的下一个终端
下一代终端

![Docker image](https://github.com/dushixiang/next-terminal/workflows/Docker%20image/badge.svg?branch=master)

Expand Down Expand Up @@ -29,6 +29,10 @@ https://next-terminal.typesafe.cn/

test/test

## 协议与条款

如您需要在企业网络中使用 next-terminal,建议先征求 IT 管理员的同意。下载、使用或分发 next-terminal 前,您必须同意 [协议](./LICENSE) 条款与限制。本项目不提供任何担保,亦不承担任何责任。

## 快速安装

- [使用docker安装](docs/install-docker.md)
Expand Down
43 changes: 43 additions & 0 deletions main.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import (
nested "github.com/antonfisher/nested-logrus-formatter"
"github.com/labstack/gommon/log"
"github.com/patrickmn/go-cache"
"github.com/robfig/cron/v3"
"github.com/sirupsen/logrus"
"gorm.io/driver/mysql"
"gorm.io/driver/sqlite"
Expand Down Expand Up @@ -155,6 +156,9 @@ func Run() error {
if err := global.DB.AutoMigrate(&model.Num{}); err != nil {
return err
}
if err := global.DB.AutoMigrate(&model.Job{}); err != nil {
return err
}

if len(model.FindAllTemp()) == 0 {
for i := 0; i <= 30; i++ {
Expand All @@ -174,6 +178,45 @@ func Run() error {
}
})
global.Store = global.NewStore()
global.Cron = cron.New(cron.WithSeconds()) //精确到秒

jobs, err := model.FindJobByFunc(model.FuncCheckAssetStatusJob)
if err != nil {
return err
}
if jobs == nil || len(jobs) == 0 {
job := model.Job{
ID: utils.UUID(),
Name: "资产状态检测",
Func: model.FuncCheckAssetStatusJob,
Cron: "0 0 0/1 * * ?",
Status: model.JobStatusRunning,
Created: utils.NowJsonTime(),
Updated: utils.NowJsonTime(),
}
if err := model.CreateNewJob(&job); err != nil {
return err
}
}

jobs, err = model.FindJobByFunc(model.FuncDelTimeoutSessionJob)
if err != nil {
return err
}
if jobs == nil || len(jobs) == 0 {
job := model.Job{
ID: utils.UUID(),
Name: "超时会话检测",
Func: model.FuncDelTimeoutSessionJob,
Cron: "0 0 0 * * ?",
Status: model.JobStatusRunning,
Created: utils.NowJsonTime(),
Updated: utils.NowJsonTime(),
}
if err := model.CreateNewJob(&job); err != nil {
return err
}
}

loginLogs, err := model.FindAliveLoginLogs()
if err != nil {
Expand Down
86 changes: 86 additions & 0 deletions pkg/api/job.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,86 @@
package api

import (
"github.com/labstack/echo/v4"
"next-terminal/pkg/model"
"strconv"
"strings"
)

func JobCreateEndpoint(c echo.Context) error {
var item model.Job
if err := c.Bind(&item); err != nil {
return err
}

if err := model.CreateNewJob(&item); err != nil {
return err
}
return Success(c, "")
}

func JobPagingEndpoint(c echo.Context) error {
pageIndex, _ := strconv.Atoi(c.QueryParam("pageIndex"))
pageSize, _ := strconv.Atoi(c.QueryParam("pageSize"))
name := c.QueryParam("name")
status := c.QueryParam("status")

items, total, err := model.FindPageJob(pageIndex, pageSize, name, status)
if err != nil {
return err
}

return Success(c, H{
"total": total,
"items": items,
})
}

func JobUpdateEndpoint(c echo.Context) error {
id := c.Param("id")

var item model.Job
if err := c.Bind(&item); err != nil {
return err
}

if err := model.UpdateJobById(&item, id); err != nil {
return err
}

return Success(c, nil)
}

func JobChangeStatusEndpoint(c echo.Context) error {
id := c.Param("id")
status := c.QueryParam("status")
if err := model.ChangeJobStatusById(id, status); err != nil {
return err
}
return Success(c, "")
}

func JobDeleteEndpoint(c echo.Context) error {
ids := c.Param("id")

split := strings.Split(ids, ",")
for i := range split {
jobId := split[i]
if err := model.DeleteJobById(jobId); err != nil {
return err
}
}

return Success(c, nil)
}

func JobGetEndpoint(c echo.Context) error {
id := c.Param("id")

item, err := model.FindJobById(id)
if err != nil {
return err
}

return Success(c, item)
}
10 changes: 10 additions & 0 deletions pkg/api/routes.go
Original file line number Diff line number Diff line change
Expand Up @@ -144,6 +144,16 @@ func SetupRoutes() *echo.Echo {
e.GET("/overview/counter", OverviewCounterEndPoint)
e.GET("/overview/sessions", OverviewSessionPoint)

jobs := e.Group("/jobs", Admin)
{
jobs.POST("", JobCreateEndpoint)
jobs.GET("/paging", JobPagingEndpoint)
jobs.PUT("/:id", JobUpdateEndpoint)
jobs.POST("/:id/change-status", JobChangeStatusEndpoint)
jobs.DELETE("/:id", JobDeleteEndpoint)
jobs.GET("/:id", JobGetEndpoint)
}

return e
}

Expand Down
3 changes: 3 additions & 0 deletions pkg/global/global.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package global

import (
"github.com/patrickmn/go-cache"
"github.com/robfig/cron/v3"
"gorm.io/gorm"
"next-terminal/pkg/config"
)
Expand All @@ -13,3 +14,5 @@ var Cache *cache.Cache
var Config *config.Config

var Store *TunStore

var Cron *cron.Cron
74 changes: 8 additions & 66 deletions pkg/handle/runner.go
Original file line number Diff line number Diff line change
@@ -1,81 +1,23 @@
package handle

import (
"github.com/robfig/cron/v3"
"github.com/sirupsen/logrus"
"log"
"next-terminal/pkg/global"
"next-terminal/pkg/guacd"
"next-terminal/pkg/model"
"next-terminal/pkg/utils"
"os"
"strconv"
"time"
)

func RunTicker() {

c := cron.New(cron.WithSeconds()) //精确到秒
// 每隔一小时删除一次未使用的会话信息
_, _ = global.Cron.AddJob("0 0 0/1 * * ?", model.DelUnUsedSessionJob{})
// 每隔一小时检测一次资产状态
//_, _ = global.Cron.AddJob("0 0 0/1 * * ?", model.CheckAssetStatusJob{})
// 每日凌晨删除超过时长限制的会话
//_, _ = global.Cron.AddJob("0 0 0 * * ?", model.DelTimeoutSessionJob{})

_, _ = c.AddFunc("0 0 0/1 * * ?", func() {
// 定时任务,每隔一小时删除一次未使用的会话信息
sessions, _ := model.FindSessionByStatusIn([]string{model.NoConnect, model.Connecting})
if sessions != nil && len(sessions) > 0 {
now := time.Now()
for i := range sessions {
if now.Sub(sessions[i].ConnectedTime.Time) > time.Hour*1 {
_ = model.DeleteSessionById(sessions[i].ID)
s := sessions[i].Username + "@" + sessions[i].IP + ":" + strconv.Itoa(sessions[i].Port)
logrus.Infof("会话「%v」ID「%v」超过1小时未打开,已删除。", s, sessions[i].ID)
}
}
}
// 每隔一小时检测一次资产是否存活
assets, _ := model.FindAllAsset()
if assets != nil && len(assets) > 0 {
for i := range assets {
asset := assets[i]
active := utils.Tcping(asset.IP, asset.Port)
model.UpdateAssetActiveById(active, asset.ID)
logrus.Infof("资产「%v」ID「%v」存活状态检测完成,存活「%v」。", asset.Name, asset.ID, active)
}
}
})

_, err := c.AddFunc("0 0 0 * * ?", func() {
// 定时任务 每日凌晨检查超过时长限制的会话
property, err := model.FindPropertyByName("session-saved-limit")
if err != nil {
return
}
if property.Value == "" || property.Value == "-" {
return
}
limit, err := strconv.Atoi(property.Value)
if err != nil {
return
}
sessions, err := model.FindOutTimeSessions(limit)
if err != nil {
return
}

if sessions != nil && len(sessions) > 0 {
var sessionIds []string
for i := range sessions {
sessionIds = append(sessionIds, sessions[i].ID)
}
err := model.DeleteSessionByIds(sessionIds)
if err != nil {
logrus.Errorf("删除离线会话失败 %v", err)
}
}
})

if err != nil {
log.Fatal(err)
}

c.Start()
global.Cron.Start()
}

func RunDataFix() {
Expand Down
Loading

0 comments on commit 1b8ecef

Please sign in to comment.