代码格式化

pull/21/merge
ouqiang 2017-04-02 10:19:52 +08:00
parent 7e5b84c320
commit 16582a5775
18 changed files with 272 additions and 286 deletions

View File

@ -1,17 +1,18 @@
package cmd package cmd
import ( import (
"github.com/urfave/cli" "github.com/go-macaron/csrf"
"gopkg.in/macaron.v1"
"github.com/go-macaron/gzip" "github.com/go-macaron/gzip"
"github.com/go-macaron/session" "github.com/go-macaron/session"
"github.com/go-macaron/csrf"
"github.com/ouqiang/cron-scheduler/modules/app" "github.com/ouqiang/cron-scheduler/modules/app"
"github.com/ouqiang/cron-scheduler/routers" "github.com/ouqiang/cron-scheduler/routers"
"github.com/urfave/cli"
"gopkg.in/macaron.v1"
) )
// web服务器默认端口 // web服务器默认端口
const DefaultPort = 5920 const DefaultPort = 5920
// 静态文件目录 // 静态文件目录
const StaticDir = "public" const StaticDir = "public"
@ -64,7 +65,7 @@ func registerMiddleware(m *macaron.Macaron) {
// 解析端口 // 解析端口
func parsePort(ctx *cli.Context) int { func parsePort(ctx *cli.Context) int {
var port int var port int
if (ctx.IsSet("port")) { if ctx.IsSet("port") {
port = ctx.Int("port") port = ctx.Int("port")
} }
if port <= 0 || port >= 65535 { if port <= 0 || port >= 65535 {

View File

@ -1,6 +1,5 @@
package models package models
// 主机 // 主机
type Host struct { type Host struct {
Id int16 `xorm:"smallint pk autoincr"` Id int16 `xorm:"smallint pk autoincr"`
@ -15,7 +14,7 @@ type Host struct {
PageSize int `xorm:"-"` PageSize int `xorm:"-"`
} }
type LoginType int8; type LoginType int8
const ( const (
PublicKey = 1 PublicKey = 1
@ -23,7 +22,7 @@ const (
) )
// 新增 // 新增
func(host *Host) Create() (insertId int16, err error) { func (host *Host) Create() (insertId int16, err error) {
_, err = Db.Insert(host) _, err = Db.Insert(host)
if err == nil { if err == nil {
insertId = host.Id insertId = host.Id
@ -33,16 +32,16 @@ func(host *Host) Create() (insertId int16, err error) {
} }
// 更新 // 更新
func(host *Host) Update(id int, data CommonMap) (int64, error) { func (host *Host) Update(id int, data CommonMap) (int64, error) {
return Db.Table(host).ID(id).Update(data) return Db.Table(host).ID(id).Update(data)
} }
// 删除 // 删除
func(host *Host) Delete(id int) (int64, error) { func (host *Host) Delete(id int) (int64, error) {
return Db.Id(id).Delete(host) return Db.Id(id).Delete(host)
} }
func(host *Host) List() ([]Host, error) { func (host *Host) List() ([]Host, error) {
host.parsePageAndPageSize() host.parsePageAndPageSize()
list := make([]Host, 0) list := make([]Host, 0)
err := Db.Desc("id").Limit(host.PageSize, host.pageLimitOffset()).Find(&list) err := Db.Desc("id").Limit(host.PageSize, host.pageLimitOffset()).Find(&list)
@ -50,19 +49,19 @@ func(host *Host) List() ([]Host, error) {
return list, err return list, err
} }
func(host *Host) Total() (int64, error) { func (host *Host) Total() (int64, error) {
return Db.Count(host) return Db.Count(host)
} }
func(host *Host) parsePageAndPageSize() { func (host *Host) parsePageAndPageSize() {
if (host.Page <= 0) { if host.Page <= 0 {
host.Page = Page host.Page = Page
} }
if (host.PageSize >= 0 || host.PageSize > MaxPageSize) { if host.PageSize >= 0 || host.PageSize > MaxPageSize {
host.PageSize = PageSize host.PageSize = PageSize
} }
} }
func(host *Host) pageLimitOffset() int { func (host *Host) pageLimitOffset() int {
return (host.Page - 1) * host.PageSize return (host.Page - 1) * host.PageSize
} }

View File

@ -4,16 +4,16 @@ import "errors"
// 创建数据库表 // 创建数据库表
type Migration struct {} type Migration struct{}
func(migration *Migration) Exec(dbName string) error { func (migration *Migration) Exec(dbName string) error {
if !isDatabaseExist(dbName) { if !isDatabaseExist(dbName) {
return errors.New("数据库不存在") return errors.New("数据库不存在")
} }
tables := []interface{}{ tables := []interface{}{
&User{}, &Task{}, &TaskLog{},&Host{}, &User{}, &Task{}, &TaskLog{}, &Host{},
} }
for _, table := range(tables) { for _, table := range tables {
err := Db.Sync2(table) err := Db.Sync2(table)
if err != nil { if err != nil {
return err return err

View File

@ -1,11 +1,11 @@
package models package models
import ( import (
"github.com/go-xorm/xorm"
"fmt" "fmt"
"github.com/ouqiang/cron-scheduler/modules/setting"
"github.com/go-xorm/core"
_ "github.com/go-sql-driver/mysql" _ "github.com/go-sql-driver/mysql"
"github.com/go-xorm/core"
"github.com/go-xorm/xorm"
"github.com/ouqiang/cron-scheduler/modules/setting"
"gopkg.in/macaron.v1" "gopkg.in/macaron.v1"
"strings" "strings"
) )
@ -30,7 +30,7 @@ const (
) )
// 创建Db // 创建Db
func CreateDb(configFile string) *xorm.Engine{ func CreateDb(configFile string) *xorm.Engine {
config := getDbConfig(configFile) config := getDbConfig(configFile)
dsn := getDbEngineDSN(config["engine"], config) dsn := getDbEngineDSN(config["engine"], config)
engine, err := xorm.NewEngine(config["engine"], dsn) engine, err := xorm.NewEngine(config["engine"], dsn)
@ -70,8 +70,8 @@ func getDbEngineDSN(engine string, config map[string]string) string {
} }
// 获取数据库配置 // 获取数据库配置
func getDbConfig(configFile string) (map[string]string){ func getDbConfig(configFile string) map[string]string {
config,err := setting.Read(configFile) config, err := setting.Read(configFile)
if err != nil { if err != nil {
panic(err) panic(err)
} }

View File

@ -39,7 +39,7 @@ type Task struct {
} }
// 新增 // 新增
func(task *Task) Create() (insertId int, err error) { func (task *Task) Create() (insertId int, err error) {
task.Status = Enabled task.Status = Enabled
_, err = Db.Insert(task) _, err = Db.Insert(task)
@ -51,26 +51,26 @@ func(task *Task) Create() (insertId int, err error) {
} }
// 更新 // 更新
func(task *Task) Update(id int, data CommonMap) (int64, error) { func (task *Task) Update(id int, data CommonMap) (int64, error) {
return Db.Table(task).ID(id).Update(data) return Db.Table(task).ID(id).Update(data)
} }
// 删除 // 删除
func(task *Task) Delete(id int) (int64, error) { func (task *Task) Delete(id int) (int64, error) {
return Db.Id(id).Delete(task) return Db.Id(id).Delete(task)
} }
// 禁用 // 禁用
func(task *Task) Disable(id int) (int64, error) { func (task *Task) Disable(id int) (int64, error) {
return task.Update(id, CommonMap{"status": Disabled}) return task.Update(id, CommonMap{"status": Disabled})
} }
// 激活 // 激活
func(task *Task) Enable(id int) (int64, error) { func (task *Task) Enable(id int) (int64, error) {
return task.Update(id, CommonMap{"status": Enabled}) return task.Update(id, CommonMap{"status": Enabled})
} }
func(task *Task) ActiveList() ([]Task, error) { func (task *Task) ActiveList() ([]Task, error) {
task.parsePageAndPageSize() task.parsePageAndPageSize()
list := make([]Task, 0) list := make([]Task, 0)
err := Db.Where("status = ?", Enabled).Desc("id").Find(&list) err := Db.Where("status = ?", Enabled).Desc("id").Find(&list)
@ -78,7 +78,7 @@ func(task *Task) ActiveList() ([]Task, error) {
return list, err return list, err
} }
func(task *Task) List() ([]Task, error) { func (task *Task) List() ([]Task, error) {
task.parsePageAndPageSize() task.parsePageAndPageSize()
list := make([]Task, 0) list := make([]Task, 0)
err := Db.Desc("id").Limit(task.PageSize, task.pageLimitOffset()).Find(&list) err := Db.Desc("id").Limit(task.PageSize, task.pageLimitOffset()).Find(&list)
@ -86,19 +86,19 @@ func(task *Task) List() ([]Task, error) {
return list, err return list, err
} }
func(taskLog *TaskLog) Total() (int64, error) { func (taskLog *TaskLog) Total() (int64, error) {
return Db.Count(taskLog) return Db.Count(taskLog)
} }
func(taskLog *TaskLog) parsePageAndPageSize() { func (taskLog *TaskLog) parsePageAndPageSize() {
if (taskLog.Page <= 0) { if taskLog.Page <= 0 {
taskLog.Page = Page taskLog.Page = Page
} }
if (taskLog.PageSize >= 0 || taskLog.PageSize > MaxPageSize) { if taskLog.PageSize >= 0 || taskLog.PageSize > MaxPageSize {
taskLog.PageSize = PageSize taskLog.PageSize = PageSize
} }
} }
func(taskLog *TaskLog) pageLimitOffset() int { func (taskLog *TaskLog) pageLimitOffset() int {
return (taskLog.Page - 1) * taskLog.PageSize return (taskLog.Page - 1) * taskLog.PageSize
} }

View File

@ -5,7 +5,7 @@ import (
) )
// 任务执行日志 // 任务执行日志
type TaskLog struct{ type TaskLog struct {
Id int `xorm:"int pk autoincr"` Id int `xorm:"int pk autoincr"`
TaskId int `xorm:"int not null"` // 任务ID TaskId int `xorm:"int not null"` // 任务ID
StartTime time.Time `xorm:"datetime created"` // 开始执行时间 StartTime time.Time `xorm:"datetime created"` // 开始执行时间
@ -16,7 +16,7 @@ type TaskLog struct{
PageSize int `xorm:"-"` PageSize int `xorm:"-"`
} }
func(taskLog *TaskLog) Create() (insertId int, err error) { func (taskLog *TaskLog) Create() (insertId int, err error) {
taskLog.Status = Running taskLog.Status = Running
_, err = Db.Insert(taskLog) _, err = Db.Insert(taskLog)
@ -28,15 +28,15 @@ func(taskLog *TaskLog) Create() (insertId int, err error) {
} }
// 更新 // 更新
func(taskLog *TaskLog) Update(id int, data CommonMap) (int64, error) { func (taskLog *TaskLog) Update(id int, data CommonMap) (int64, error) {
return Db.Table(taskLog).ID(id).Update(data) return Db.Table(taskLog).ID(id).Update(data)
} }
func(taskLog *TaskLog) setStatus(id int ,status Status) (int64, error) { func (taskLog *TaskLog) setStatus(id int, status Status) (int64, error) {
return taskLog.Update(id, CommonMap{"status": status}) return taskLog.Update(id, CommonMap{"status": status})
} }
func(taskLog *TaskLog) List() ([]TaskLog, error) { func (taskLog *TaskLog) List() ([]TaskLog, error) {
taskLog.parsePageAndPageSize() taskLog.parsePageAndPageSize()
list := make([]TaskLog, 0) list := make([]TaskLog, 0)
err := Db.Desc("id").Limit(taskLog.PageSize, taskLog.pageLimitOffset()).Find(&list) err := Db.Desc("id").Limit(taskLog.PageSize, taskLog.pageLimitOffset()).Find(&list)
@ -44,19 +44,19 @@ func(taskLog *TaskLog) List() ([]TaskLog, error) {
return list, err return list, err
} }
func(task *Task) Total() (int64, error) { func (task *Task) Total() (int64, error) {
return Db.Count(task) return Db.Count(task)
} }
func(task *Task) parsePageAndPageSize() { func (task *Task) parsePageAndPageSize() {
if (task.Page <= 0) { if task.Page <= 0 {
task.Page = Page task.Page = Page
} }
if (task.PageSize >= 0 || task.PageSize > MaxPageSize) { if task.PageSize >= 0 || task.PageSize > MaxPageSize {
task.PageSize = PageSize task.PageSize = PageSize
} }
} }
func(task *Task) pageLimitOffset() int { func (task *Task) pageLimitOffset() int {
return (task.Page - 1) * task.PageSize return (task.Page - 1) * task.PageSize
} }

View File

@ -1,11 +1,11 @@
package models package models
import ( import (
"time"
"github.com/ouqiang/cron-scheduler/modules/utils" "github.com/ouqiang/cron-scheduler/modules/utils"
"time"
) )
const PasswordSaltLength = 6; const PasswordSaltLength = 6
// 用户model // 用户model
type User struct { type User struct {
@ -24,7 +24,7 @@ type User struct {
} }
// 新增 // 新增
func(user *User) Create() (insertId int, err error) { func (user *User) Create() (insertId int, err error) {
user.Status = Enabled user.Status = Enabled
user.Salt = user.generateSalt() user.Salt = user.generateSalt()
user.Password = user.encryptPassword(user.Password, user.Salt) user.Password = user.encryptPassword(user.Password, user.Salt)
@ -38,34 +38,34 @@ func(user *User) Create() (insertId int, err error) {
} }
// 更新 // 更新
func(user *User) Update(id int, data CommonMap) (int64, error) { func (user *User) Update(id int, data CommonMap) (int64, error) {
return Db.Table(user).ID(id).Update(data) return Db.Table(user).ID(id).Update(data)
} }
// 删除 // 删除
func(user *User) Delete(id int) (int64, error) { func (user *User) Delete(id int) (int64, error) {
return Db.Id(id).Delete(user) return Db.Id(id).Delete(user)
} }
// 禁用 // 禁用
func(user *User) Disable(id int) (int64, error) { func (user *User) Disable(id int) (int64, error) {
return user.Update(id, CommonMap{"status": Disabled}) return user.Update(id, CommonMap{"status": Disabled})
} }
// 激活 // 激活
func(user *User) Enable(id int) (int64, error) { func (user *User) Enable(id int) (int64, error) {
return user.Update(id, CommonMap{"status": Enabled}) return user.Update(id, CommonMap{"status": Enabled})
} }
// 验证用户名和密码 // 验证用户名和密码
func(user *User) Match(username, password string) bool { func (user *User) Match(username, password string) bool {
where := "(name = ? OR email = ?)" where := "(name = ? OR email = ?)"
_, err := Db.Where(where, username, username).Get(user) _, err := Db.Where(where, username, username).Get(user)
if err != nil { if err != nil {
return false return false
} }
hashPassword := user.encryptPassword(password, user.Salt) hashPassword := user.encryptPassword(password, user.Salt)
if (hashPassword != user.Password) { if hashPassword != user.Password {
return false return false
} }
@ -73,17 +73,16 @@ func(user *User) Match(username, password string) bool {
} }
// 用户名是否存在 // 用户名是否存在
func(user *User) UsernameExists(username string ) (int64, error) { func (user *User) UsernameExists(username string) (int64, error) {
return Db.Where("name = ?", username).Count(user) return Db.Where("name = ?", username).Count(user)
} }
// 邮箱地址是否存在 // 邮箱地址是否存在
func(user *User) EmailExists(email string) (int64, error) { func (user *User) EmailExists(email string) (int64, error) {
return Db.Where("email = ?", email).Count(user) return Db.Where("email = ?", email).Count(user)
} }
func (user *User) List() ([]User, error) {
func(user *User) List() ([]User, error) {
user.parsePageAndPageSize() user.parsePageAndPageSize()
list := make([]User, 0) list := make([]User, 0)
err := Db.Desc("id").Limit(user.PageSize, user.pageLimitOffset()).Find(&list) err := Db.Desc("id").Limit(user.PageSize, user.pageLimitOffset()).Find(&list)
@ -91,29 +90,29 @@ func(user *User) List() ([]User, error) {
return list, err return list, err
} }
func(user *User) Total() (int64, error) { func (user *User) Total() (int64, error) {
return Db.Count(user) return Db.Count(user)
} }
func(user *User) parsePageAndPageSize() { func (user *User) parsePageAndPageSize() {
if (user.Page <= 0) { if user.Page <= 0 {
user.Page = Page user.Page = Page
} }
if (user.PageSize >= 0 || user.PageSize > MaxPageSize) { if user.PageSize >= 0 || user.PageSize > MaxPageSize {
user.PageSize = PageSize user.PageSize = PageSize
} }
} }
func(user *User) pageLimitOffset() int { func (user *User) pageLimitOffset() int {
return (user.Page - 1) * user.PageSize return (user.Page - 1) * user.PageSize
} }
// 密码加密 // 密码加密
func(user *User) encryptPassword(password, salt string) string { func (user *User) encryptPassword(password, salt string) string {
return utils.Md5(password + salt) return utils.Md5(password + salt)
} }
// 生成密码盐值 // 生成密码盐值
func(user *User) generateSalt() string { func (user *User) generateSalt() string {
return utils.RandString(PasswordSaltLength) return utils.RandString(PasswordSaltLength)
} }

View File

@ -7,17 +7,15 @@ import (
"github.com/ouqiang/cron-scheduler/modules/utils" "github.com/ouqiang/cron-scheduler/modules/utils"
) )
/** /**
* ad-hoc * ad-hoc
* hosts * hosts
* hostFile * hostFile
* module * module
* args module * args module
*/ */
func ExecCommand(hosts string, hostFile string, module string, args... string) (output string, err error) { func ExecCommand(hosts string, hostFile string, module string, args ...string) (output string, err error) {
if hosts== "" || hostFile == "" || module == "" { if hosts == "" || hostFile == "" || module == "" {
err = errors.New("参数不完整") err = errors.New("参数不完整")
return return
} }
@ -31,12 +29,12 @@ func ExecCommand(hosts string, hostFile string, module string, args... string) (
} }
// 执行shell命令 // 执行shell命令
func Shell(hosts string, hostFile string, args... string) (output string, err error) { func Shell(hosts string, hostFile string, args ...string) (output string, err error) {
return ExecCommand(hosts, hostFile, "shell", args...) return ExecCommand(hosts, hostFile, "shell", args...)
} }
// 复制本地脚本到远程执行 // 复制本地脚本到远程执行
func Script(hosts string, hostFile string, args... string) (output string, err error) { func Script(hosts string, hostFile string, args ...string) (output string, err error) {
return ExecCommand(hosts, hostFile, "script", args...) return ExecCommand(hosts, hostFile, "script", args...)
} }

View File

@ -1,12 +1,12 @@
package ansible package ansible
import ( import (
"github.com/ouqiang/cron-scheduler/models"
"sync"
"io/ioutil"
"bytes" "bytes"
"strconv" "github.com/ouqiang/cron-scheduler/models"
"github.com/ouqiang/cron-scheduler/modules/utils" "github.com/ouqiang/cron-scheduler/modules/utils"
"io/ioutil"
"strconv"
"sync"
) )
// 主机名 // 主机名
@ -25,16 +25,15 @@ func NewHosts(hostFilename string) *Hosts {
} }
// 获取hosts文件名 // 获取hosts文件名
func(h *Hosts) GetFilename() string { func (h *Hosts) GetFilename() string {
h.RLock() h.RLock()
defer h.RUnlock() defer h.RUnlock()
return h.filename return h.filename
} }
// 写入hosts // 写入hosts
func(h *Hosts) Write() { func (h *Hosts) Write() {
host := new(models.Host) host := new(models.Host)
hostModels, err := host.List() hostModels, err := host.List()
if err != nil { if err != nil {
@ -46,7 +45,7 @@ func(h *Hosts) Write() {
return return
} }
buffer := bytes.Buffer{} buffer := bytes.Buffer{}
for _, hostModel := range(hostModels) { for _, hostModel := range hostModels {
buffer.WriteString(strconv.Itoa(int(hostModel.Id))) buffer.WriteString(strconv.Itoa(int(hostModel.Id)))
buffer.WriteString(" ansible_ssh_host=") buffer.WriteString(" ansible_ssh_host=")
buffer.WriteString(hostModel.Name) buffer.WriteString(hostModel.Name)
@ -54,7 +53,7 @@ func(h *Hosts) Write() {
buffer.WriteString(strconv.Itoa(hostModel.Port)) buffer.WriteString(strconv.Itoa(hostModel.Port))
buffer.WriteString(" ansible_ssh_user=") buffer.WriteString(" ansible_ssh_user=")
buffer.WriteString(hostModel.Username) buffer.WriteString(hostModel.Username)
if (hostModel.LoginType != models.PublicKey && hostModel.Password != "") { if hostModel.LoginType != models.PublicKey && hostModel.Password != "" {
buffer.WriteString(" ansible_ssh_pass=") buffer.WriteString(" ansible_ssh_pass=")
buffer.WriteString(hostModel.Password) buffer.WriteString(hostModel.Password)
} }
@ -66,5 +65,3 @@ func(h *Hosts) Write() {
return return
} }

View File

@ -4,10 +4,10 @@ import (
"os" "os"
"runtime" "runtime"
"github.com/ouqiang/cron-scheduler/modules/utils"
"github.com/ouqiang/cron-scheduler/modules/crontask"
"github.com/ouqiang/cron-scheduler/models" "github.com/ouqiang/cron-scheduler/models"
"github.com/ouqiang/cron-scheduler/modules/ansible" "github.com/ouqiang/cron-scheduler/modules/ansible"
"github.com/ouqiang/cron-scheduler/modules/crontask"
"github.com/ouqiang/cron-scheduler/modules/utils"
"github.com/ouqiang/cron-scheduler/service" "github.com/ouqiang/cron-scheduler/service"
) )
@ -79,7 +79,6 @@ func CreateInstallLock() error {
return err return err
} }
// 初始化资源 // 初始化资源
func InitResource() { func InitResource() {
// 初始化ansible Hosts // 初始化ansible Hosts
@ -96,8 +95,8 @@ func InitDb() {
} }
// 检测目录是否存在 // 检测目录是否存在
func checkDirExists(path... string) { func checkDirExists(path ...string) {
for _, value := range(path) { for _, value := range path {
_, err := os.Stat(value) _, err := os.Stat(value)
if os.IsNotExist(err) { if os.IsNotExist(err) {
panic(value + "目录不存在") panic(value + "目录不存在")

View File

@ -1,10 +1,10 @@
package crontask package crontask
import ( import (
"github.com/robfig/cron"
"errors" "errors"
"sync" "github.com/robfig/cron"
"strings" "strings"
"sync"
) )
var DefaultCronTask *CronTask var DefaultCronTask *CronTask
@ -17,7 +17,7 @@ type CronTask struct {
} }
func NewCronTask() *CronTask { func NewCronTask() *CronTask {
return &CronTask { return &CronTask{
sync.RWMutex{}, sync.RWMutex{},
make(CronMap), make(CronMap),
} }
@ -26,7 +26,7 @@ func NewCronTask() *CronTask {
// 新增定时任务,如果name存在则添加失败 // 新增定时任务,如果name存在则添加失败
// name 任务名称 // name 任务名称
// spec crontab时间格式定义 可定义多个时间\n分隔 // spec crontab时间格式定义 可定义多个时间\n分隔
func(cronTask *CronTask) Add(name string, spec string, cmd cron.FuncJob ) (err error) { func (cronTask *CronTask) Add(name string, spec string, cmd cron.FuncJob) (err error) {
if name == "" || spec == "" || cmd == nil { if name == "" || spec == "" || cmd == nil {
return errors.New("参数不完整") return errors.New("参数不完整")
} }
@ -39,13 +39,13 @@ func(cronTask *CronTask) Add(name string, spec string, cmd cron.FuncJob ) (err e
defer cronTask.Unlock() defer cronTask.Unlock()
cronTask.tasks[name] = cron.New() cronTask.tasks[name] = cron.New()
specs := strings.Split(spec, "|||") specs := strings.Split(spec, "|||")
for _, item := range(specs) { for _, item := range specs {
_, err = cron.Parse(item) _, err = cron.Parse(item)
if err != nil { if err != nil {
return err return err
} }
} }
for _, item := range(specs) { for _, item := range specs {
err = cronTask.tasks[name].AddFunc(item, cmd) err = cronTask.tasks[name].AddFunc(item, cmd)
} }
cronTask.tasks[name].Start() cronTask.tasks[name].Start()
@ -54,7 +54,7 @@ func(cronTask *CronTask) Add(name string, spec string, cmd cron.FuncJob ) (err e
} }
// 任务不存在则新增,任务已存在则删除后新增 // 任务不存在则新增,任务已存在则删除后新增
func(cronTask *CronTask) AddOrReplace(name string, spec string, cmd cron.FuncJob) error { func (cronTask *CronTask) AddOrReplace(name string, spec string, cmd cron.FuncJob) error {
if cronTask.IsExist(name) { if cronTask.IsExist(name) {
cronTask.Delete(name) cronTask.Delete(name)
} }
@ -62,9 +62,8 @@ func(cronTask *CronTask) AddOrReplace(name string, spec string, cmd cron.FuncJob
return cronTask.Add(name, spec, cmd) return cronTask.Add(name, spec, cmd)
} }
// 判断任务是否存在 // 判断任务是否存在
func(cronTask *CronTask) IsExist(name string) bool { func (cronTask *CronTask) IsExist(name string) bool {
cronTask.RLock() cronTask.RLock()
defer cronTask.RUnlock() defer cronTask.RUnlock()
_, ok := cronTask.tasks[name] _, ok := cronTask.tasks[name]
@ -73,14 +72,14 @@ func(cronTask *CronTask) IsExist(name string) bool {
} }
// 停止任务 // 停止任务
func(cronTask *CronTask) Stop(name string) { func (cronTask *CronTask) Stop(name string) {
if cronTask.IsExist(name) { if cronTask.IsExist(name) {
cronTask.tasks[name].Stop() cronTask.tasks[name].Stop()
} }
} }
// 删除任务 // 删除任务
func(cronTask *CronTask) Delete(name string) { func (cronTask *CronTask) Delete(name string) {
cronTask.Stop(name) cronTask.Stop(name)
cronTask.Lock() cronTask.Lock()
defer cronTask.Unlock() defer cronTask.Unlock()

View File

@ -1,8 +1,8 @@
package setting package setting
import ( import (
"gopkg.in/ini.v1"
"errors" "errors"
"gopkg.in/ini.v1"
) )
// 读取配置 // 读取配置
@ -15,15 +15,14 @@ func Read(filename string) (config *ini.File, err error) {
return return
} }
// 写入配置 // 写入配置
func Write(config map[string]map[string]string, filename string) (error) { func Write(config map[string]map[string]string, filename string) error {
if len(config) == 0 { if len(config) == 0 {
return errors.New("参数不能为空") return errors.New("参数不能为空")
} }
file := ini.Empty() file := ini.Empty()
for sectionName, items := range(config) { for sectionName, items := range config {
if sectionName == "" { if sectionName == "" {
return errors.New("节名称不能为空") return errors.New("节名称不能为空")
} }
@ -31,7 +30,7 @@ func Write(config map[string]map[string]string, filename string) (error) {
if err != nil { if err != nil {
return err return err
} }
for key, value := range(items) { for key, value := range items {
_, err = section.NewKey(key, value) _, err = section.NewKey(key, value)
if err != nil { if err != nil {
return err return err

View File

@ -10,20 +10,20 @@ type response struct {
Data interface{} `json:"data"` // 数据 Data interface{} `json:"data"` // 数据
} }
type Json struct {} type Json struct{}
const ResponseSuccess = 0; const ResponseSuccess = 0
const ResponseFailure = 1; const ResponseFailure = 1
func(j *Json) Success(message string, data interface{}) string { func (j *Json) Success(message string, data interface{}) string {
return j.response(ResponseSuccess, message, data) return j.response(ResponseSuccess, message, data)
} }
func(j *Json) Failure(code int, message string) string { func (j *Json) Failure(code int, message string) string {
return j.response(code, message, nil) return j.response(code, message, nil)
} }
func(j *Json) response(code int, message string, data interface{}) (string) { func (j *Json) response(code int, message string, data interface{}) string {
resp := response{ resp := response{
Code: code, Code: code,
Message: message, Message: message,
@ -37,4 +37,3 @@ func(j *Json) response(code int, message string, data interface{}) (string) {
return string(result) return string(result)
} }

View File

@ -1,17 +1,16 @@
package utils package utils
import ( import (
"os/exec"
"math/rand"
"time"
"crypto/md5" "crypto/md5"
"encoding/hex" "encoding/hex"
"log" "log"
"math/rand"
"os/exec"
"time"
) )
// 执行shell命令 // 执行shell命令
func ExecShell(command string, args... string) (string, error) { func ExecShell(command string, args ...string) (string, error) {
result, err := exec.Command(command, args...).CombinedOutput() result, err := exec.Command(command, args...).CombinedOutput()
return string(result), err return string(result), err
@ -48,6 +47,6 @@ func RandNumber(max int) int {
// 日志记录 // 日志记录
// todo 保存到哪里 文件,数据库还是elasticsearch?,暂时输出到终端 // todo 保存到哪里 文件,数据库还是elasticsearch?,暂时输出到终端
func RecordLog(v... interface{}) { func RecordLog(v ...interface{}) {
log.Println(v) log.Println(v)
} }

View File

@ -1,11 +1,11 @@
package install package install
import ( import (
"gopkg.in/macaron.v1"
"github.com/ouqiang/cron-scheduler/modules/app"
"github.com/ouqiang/cron-scheduler/modules/utils"
"github.com/ouqiang/cron-scheduler/models" "github.com/ouqiang/cron-scheduler/models"
"github.com/ouqiang/cron-scheduler/modules/app"
"github.com/ouqiang/cron-scheduler/modules/setting" "github.com/ouqiang/cron-scheduler/modules/setting"
"github.com/ouqiang/cron-scheduler/modules/utils"
"gopkg.in/macaron.v1"
"strconv" "strconv"
) )

View File

@ -1,9 +1,9 @@
package routers package routers
import ( import (
"gopkg.in/macaron.v1"
"github.com/ouqiang/cron-scheduler/routers/install"
"github.com/go-macaron/binding" "github.com/go-macaron/binding"
"github.com/ouqiang/cron-scheduler/routers/install"
"gopkg.in/macaron.v1"
) )
// 路由注册 // 路由注册
@ -19,7 +19,7 @@ func Register(m *macaron.Macaron) {
ctx.HTML(500, "error/500") ctx.HTML(500, "error/500")
}) })
// 首页 // 首页
m.Get("/", func(ctx *macaron.Context) (string) { m.Get("/", func(ctx *macaron.Context) string {
return "go home" return "go home"
}) })
// 系统安装 // 系统安装

View File

@ -1,22 +1,22 @@
package service package service
import ( import (
"fmt"
"github.com/ouqiang/cron-scheduler/models" "github.com/ouqiang/cron-scheduler/models"
"github.com/ouqiang/cron-scheduler/modules/ansible"
"github.com/ouqiang/cron-scheduler/modules/crontask"
"github.com/ouqiang/cron-scheduler/modules/utils" "github.com/ouqiang/cron-scheduler/modules/utils"
"net/http" "github.com/robfig/cron"
"io/ioutil" "io/ioutil"
"net/http"
"strconv" "strconv"
"time" "time"
"github.com/ouqiang/cron-scheduler/modules/crontask"
"github.com/robfig/cron"
"github.com/ouqiang/cron-scheduler/modules/ansible"
"fmt"
) )
type Task struct {} type Task struct{}
// 初始化任务, 从数据库取出所有任务, 添加到定时任务并运行 // 初始化任务, 从数据库取出所有任务, 添加到定时任务并运行
func(task *Task) Initialize() { func (task *Task) Initialize() {
taskModel := new(models.Task) taskModel := new(models.Task)
taskList, err := taskModel.ActiveList() taskList, err := taskModel.ActiveList()
if err != nil { if err != nil {
@ -27,15 +27,13 @@ func(task *Task) Initialize() {
utils.RecordLog("任务列表为空") utils.RecordLog("任务列表为空")
return return
} }
for _, item := range(taskList) { for _, item := range taskList {
task.Add(item) task.Add(item)
} }
} }
// 添加任务 // 添加任务
func(task *Task) Add(taskModel models.Task) { func (task *Task) Add(taskModel models.Task) {
taskFunc := createHandlerJob(taskModel) taskFunc := createHandlerJob(taskModel)
if taskFunc == nil { if taskFunc == nil {
utils.RecordLog("添加任务#不存在的任务协议编号", taskModel.Protocol) utils.RecordLog("添加任务#不存在的任务协议编号", taskModel.Protocol)
@ -49,7 +47,7 @@ func(task *Task) Add(taskModel models.Task) {
} }
} else if taskModel.Type == models.Delay { } else if taskModel.Type == models.Delay {
// 延时任务 // 延时任务
time.AfterFunc(time.Duration(taskModel.Delay) * time.Second, taskFunc) time.AfterFunc(time.Duration(taskModel.Delay)*time.Second, taskFunc)
} }
} }
@ -58,11 +56,11 @@ type Handler interface {
} }
// HTTP任务 // HTTP任务
type HTTPHandler struct {} type HTTPHandler struct{}
func(h *HTTPHandler) Run(taskModel models.Task) (result string, err error) { func (h *HTTPHandler) Run(taskModel models.Task) (result string, err error) {
client := &http.Client{} client := &http.Client{}
if (taskModel.Timeout > 0) { if taskModel.Timeout > 0 {
client.Timeout = time.Duration(taskModel.Timeout) * time.Second client.Timeout = time.Duration(taskModel.Timeout) * time.Second
} }
req, err := http.NewRequest("POST", taskModel.Command, nil) req, err := http.NewRequest("POST", taskModel.Command, nil)
@ -88,28 +86,27 @@ func(h *HTTPHandler) Run(taskModel models.Task) (result string, err error) {
utils.RecordLog("读取HTTP请求返回值失败-", err.Error()) utils.RecordLog("读取HTTP请求返回值失败-", err.Error())
} }
return string(body),err return string(body), err
} }
// SSH-command任务 // SSH-command任务
type SSHCommandHandler struct {} type SSHCommandHandler struct{}
func(ssh *SSHCommandHandler) Run(taskModel models.Task) (string, error) { func (ssh *SSHCommandHandler) Run(taskModel models.Task) (string, error) {
return execSSHHandler("shell", taskModel) return execSSHHandler("shell", taskModel)
} }
// SSH-script任务 // SSH-script任务
type SSHScriptHandler struct {} type SSHScriptHandler struct{}
func(ssh *SSHScriptHandler) Run(taskModel models.Task) (string, error) { func (ssh *SSHScriptHandler) Run(taskModel models.Task) (string, error) {
return execSSHHandler("script", taskModel) return execSSHHandler("script", taskModel)
} }
// SSH任务 // SSH任务
func execSSHHandler(module string, taskModel models.Task) (string, error) { func execSSHHandler(module string, taskModel models.Task) (string, error) {
var args []string = []string{ taskModel.Command } var args []string = []string{taskModel.Command}
if (taskModel.Timeout > 0) { if taskModel.Timeout > 0 {
// -B 异步执行超时时间, -P 轮询时间 // -B 异步执行超时时间, -P 轮询时间
args = append(args, "-B", strconv.Itoa(taskModel.Timeout), "-P", "10") args = append(args, "-B", strconv.Itoa(taskModel.Timeout), "-P", "10")
} }
@ -146,7 +143,7 @@ func updateTaskLog(taskLogId int, result string, err error) (int64, error) {
return taskLogModel.Update(taskLogId, models.CommonMap{ return taskLogModel.Update(taskLogId, models.CommonMap{
"status": status, "status": status,
"result": result, "result": result,
}); })
} }