302 lines
13 KiB
Go
302 lines
13 KiB
Go
package boot
|
||
|
||
import (
|
||
"context"
|
||
"time"
|
||
|
||
"ppgo_job/service/scheduler"
|
||
|
||
"github.com/gogf/gf/v2/frame/g"
|
||
)
|
||
|
||
func init() {
|
||
ctx := context.Background()
|
||
g.Log().Info(ctx, "PPGo_Job 正在初始化...")
|
||
|
||
autoInitDB(ctx)
|
||
scheduler.Scheduler.Init(ctx)
|
||
|
||
g.Log().Info(ctx, "PPGo_Job 初始化完成")
|
||
}
|
||
|
||
func autoInitDB(ctx context.Context) {
|
||
tables := []string{
|
||
`CREATE TABLE IF NOT EXISTS pp_uc_admin (
|
||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||
login_name TEXT NOT NULL UNIQUE, real_name TEXT NOT NULL DEFAULT '',
|
||
password TEXT NOT NULL DEFAULT '', role_ids TEXT NOT NULL DEFAULT '',
|
||
phone TEXT NOT NULL DEFAULT '', email TEXT NOT NULL DEFAULT '',
|
||
dingtalk TEXT NOT NULL DEFAULT '', wechat TEXT NOT NULL DEFAULT '',
|
||
salt TEXT NOT NULL DEFAULT '', last_login INTEGER NOT NULL DEFAULT 0,
|
||
last_ip TEXT NOT NULL DEFAULT '', status INTEGER NOT NULL DEFAULT 1,
|
||
create_id INTEGER NOT NULL DEFAULT 0, update_id INTEGER NOT NULL DEFAULT 0,
|
||
create_time INTEGER NOT NULL DEFAULT 0, update_time INTEGER NOT NULL DEFAULT 0
|
||
)`,
|
||
`CREATE TABLE IF NOT EXISTS pp_uc_auth (
|
||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||
auth_name TEXT NOT NULL DEFAULT '', auth_url TEXT NOT NULL DEFAULT '',
|
||
user_id INTEGER NOT NULL DEFAULT 0, pid INTEGER NOT NULL DEFAULT 0,
|
||
sort INTEGER NOT NULL DEFAULT 0, icon TEXT NOT NULL DEFAULT '',
|
||
is_show INTEGER NOT NULL DEFAULT 1, status INTEGER NOT NULL DEFAULT 1,
|
||
create_id INTEGER NOT NULL DEFAULT 0, update_id INTEGER NOT NULL DEFAULT 0,
|
||
create_time INTEGER NOT NULL DEFAULT 0, update_time INTEGER NOT NULL DEFAULT 0
|
||
)`,
|
||
`CREATE TABLE IF NOT EXISTS pp_uc_role (
|
||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||
role_name TEXT NOT NULL DEFAULT '', detail TEXT NOT NULL DEFAULT '',
|
||
server_group_ids TEXT NOT NULL DEFAULT '', task_group_ids TEXT NOT NULL DEFAULT '',
|
||
status INTEGER NOT NULL DEFAULT 1,
|
||
create_id INTEGER NOT NULL DEFAULT 0, update_id INTEGER NOT NULL DEFAULT 0,
|
||
create_time INTEGER NOT NULL DEFAULT 0, update_time INTEGER NOT NULL DEFAULT 0
|
||
)`,
|
||
`CREATE TABLE IF NOT EXISTS pp_uc_role_auth (
|
||
auth_id INTEGER NOT NULL, role_id INTEGER NOT NULL,
|
||
PRIMARY KEY (auth_id, role_id)
|
||
)`,
|
||
`CREATE TABLE IF NOT EXISTS pp_task_server_group (
|
||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||
group_name TEXT NOT NULL DEFAULT '', description TEXT NOT NULL DEFAULT '',
|
||
create_id INTEGER NOT NULL DEFAULT 0, update_id INTEGER NOT NULL DEFAULT 0,
|
||
create_time INTEGER NOT NULL DEFAULT 0, update_time INTEGER NOT NULL DEFAULT 0,
|
||
status INTEGER NOT NULL DEFAULT 1
|
||
)`,
|
||
`CREATE TABLE IF NOT EXISTS pp_task_server (
|
||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||
group_id INTEGER NOT NULL DEFAULT 0, connection_type INTEGER NOT NULL DEFAULT 0,
|
||
server_name TEXT NOT NULL DEFAULT '', server_account TEXT NOT NULL DEFAULT '',
|
||
server_outer_ip TEXT NOT NULL DEFAULT '', server_ip TEXT NOT NULL DEFAULT '',
|
||
port INTEGER NOT NULL DEFAULT 22, password TEXT NOT NULL DEFAULT '',
|
||
private_key_src TEXT NOT NULL DEFAULT '', public_key_src TEXT NOT NULL DEFAULT '',
|
||
type INTEGER NOT NULL DEFAULT 0, detail TEXT NOT NULL DEFAULT '',
|
||
create_time INTEGER NOT NULL DEFAULT 0, update_time INTEGER NOT NULL DEFAULT 0,
|
||
status INTEGER NOT NULL DEFAULT 1
|
||
)`,
|
||
`CREATE TABLE IF NOT EXISTS pp_task_group (
|
||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||
group_name TEXT NOT NULL DEFAULT '', description TEXT NOT NULL DEFAULT '',
|
||
create_id INTEGER NOT NULL DEFAULT 0, update_id INTEGER NOT NULL DEFAULT 0,
|
||
create_time INTEGER NOT NULL DEFAULT 0, update_time INTEGER NOT NULL DEFAULT 0,
|
||
status INTEGER NOT NULL DEFAULT 1
|
||
)`,
|
||
`CREATE TABLE IF NOT EXISTS pp_task (
|
||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||
group_id INTEGER NOT NULL DEFAULT 0, server_ids TEXT NOT NULL DEFAULT '',
|
||
server_type INTEGER NOT NULL DEFAULT 0, task_name TEXT NOT NULL DEFAULT '',
|
||
description TEXT NOT NULL DEFAULT '', cron_spec TEXT NOT NULL DEFAULT '',
|
||
concurrent INTEGER NOT NULL DEFAULT 0, task_type TEXT NOT NULL DEFAULT 'shell',
|
||
command TEXT NOT NULL DEFAULT '', url TEXT NOT NULL DEFAULT '',
|
||
method TEXT NOT NULL DEFAULT 'GET', headers TEXT NOT NULL DEFAULT '',
|
||
body TEXT NOT NULL DEFAULT '',
|
||
timeout INTEGER NOT NULL DEFAULT 0, execute_times INTEGER NOT NULL DEFAULT 0,
|
||
prev_time INTEGER NOT NULL DEFAULT 0, status INTEGER NOT NULL DEFAULT 2,
|
||
is_notify INTEGER NOT NULL DEFAULT 0, notify_type INTEGER NOT NULL DEFAULT 0,
|
||
notify_tpl_id INTEGER NOT NULL DEFAULT 0, notify_user_ids TEXT NOT NULL DEFAULT '',
|
||
create_id INTEGER NOT NULL DEFAULT 0, update_id INTEGER NOT NULL DEFAULT 0,
|
||
create_time INTEGER NOT NULL DEFAULT 0, update_time INTEGER NOT NULL DEFAULT 0
|
||
)`,
|
||
`CREATE TABLE IF NOT EXISTS pp_task_log (
|
||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||
task_id INTEGER NOT NULL DEFAULT 0, server_id INTEGER NOT NULL DEFAULT 0,
|
||
server_name TEXT NOT NULL DEFAULT '', output TEXT NOT NULL DEFAULT '',
|
||
error TEXT NOT NULL DEFAULT '', status INTEGER NOT NULL DEFAULT 0,
|
||
process_time INTEGER NOT NULL DEFAULT 0, create_time INTEGER NOT NULL DEFAULT 0
|
||
)`,
|
||
`CREATE TABLE IF NOT EXISTS pp_task_ban (
|
||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||
code TEXT NOT NULL DEFAULT '', create_time INTEGER NOT NULL DEFAULT 0,
|
||
update_time INTEGER NOT NULL DEFAULT 0, status INTEGER NOT NULL DEFAULT 1
|
||
)`,
|
||
`CREATE TABLE IF NOT EXISTS pp_notify_tpl (
|
||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||
type TEXT NOT NULL DEFAULT 'system', tpl_name TEXT NOT NULL DEFAULT '',
|
||
tpl_type INTEGER NOT NULL DEFAULT 0, title TEXT NOT NULL DEFAULT '',
|
||
content TEXT NOT NULL DEFAULT '', status INTEGER NOT NULL DEFAULT 1,
|
||
create_id INTEGER NOT NULL DEFAULT 0, update_id INTEGER NOT NULL DEFAULT 0,
|
||
create_time INTEGER NOT NULL DEFAULT 0, update_time INTEGER NOT NULL DEFAULT 0
|
||
)`,
|
||
}
|
||
|
||
for _, sql := range tables {
|
||
if _, err := g.DB().Exec(ctx, sql); err != nil {
|
||
g.Log().Warningf(ctx, "建表失败: %v", err)
|
||
}
|
||
}
|
||
|
||
// 迁移:为旧表添加新列(忽略重复添加的错误)
|
||
migrations := []string{
|
||
"ALTER TABLE pp_task ADD COLUMN task_type TEXT NOT NULL DEFAULT 'shell'",
|
||
"ALTER TABLE pp_task ADD COLUMN url TEXT NOT NULL DEFAULT ''",
|
||
"ALTER TABLE pp_task ADD COLUMN method TEXT NOT NULL DEFAULT 'GET'",
|
||
"ALTER TABLE pp_task ADD COLUMN headers TEXT NOT NULL DEFAULT ''",
|
||
"ALTER TABLE pp_task ADD COLUMN body TEXT NOT NULL DEFAULT ''",
|
||
}
|
||
for _, sql := range migrations {
|
||
if _, err := g.DB().Exec(ctx, sql); err != nil {
|
||
// 列已存在则忽略错误
|
||
}
|
||
}
|
||
|
||
var adminCount int
|
||
_ = g.DB().GetScan(ctx, &adminCount, "SELECT COUNT(*) FROM pp_uc_admin WHERE id=?", 1)
|
||
if adminCount == 0 {
|
||
_, _ = g.DB().Exec(ctx,
|
||
"INSERT INTO pp_uc_admin(id, login_name, real_name, password, role_ids, salt, status, create_time, update_time) VALUES(?,?,?,?,?,?,?,?,?)",
|
||
1, "admin", "管理员", "bc9b5718afdffe85fb13555347969ff5", "0", "abcd", 1, time.Now().Unix(), time.Now().Unix(),
|
||
)
|
||
g.Log().Info(ctx, "已创建默认管理员 admin / 123456")
|
||
}
|
||
|
||
var authCount int
|
||
_ = g.DB().GetScan(ctx, &authCount, "SELECT COUNT(*) FROM pp_uc_auth")
|
||
if authCount == 0 {
|
||
seedAuth(ctx)
|
||
}
|
||
|
||
// 种子数据:业务数据(任务分组 + 定时任务)
|
||
seedBusinessData(ctx)
|
||
}
|
||
|
||
func seedAuth(ctx context.Context) {
|
||
now := time.Now().Unix()
|
||
type row struct {
|
||
id, pid, sort, isShow int
|
||
name, url, icon string
|
||
}
|
||
rows := []row{
|
||
// pid=0 = 一级菜单
|
||
{1, 0, 1, 1, "任务管理", " ", "fa-tasks"},
|
||
{2, 0, 2, 1, "服务器管理", " ", "fa-server"},
|
||
{3, 0, 3, 1, "系统设置", " ", "fa-cog"},
|
||
{4, 0, 4, 1, "日志管理", " ", "fa-file-text"},
|
||
{5, 0, 5, 1, "权限管理", " ", "fa-lock"},
|
||
// 二级菜单(pid=父菜单的id)
|
||
{14, 1, 1, 1, "任务列表", "/task/list", "fa-tasks"},
|
||
{15, 1, 2, 1, "任务审核", "/task/audit_list", "fa-check"},
|
||
{16, 2, 1, 1, "服务器列表", "/server/list", "fa-server"},
|
||
{6, 3, 1, 1, "任务分组", "/group/list", "fa-folder"},
|
||
{7, 3, 2, 1, "资源分组", "/server_group/list", "fa-sitemap"},
|
||
{8, 3, 3, 1, "禁用命令", "/ban/list", "fa-ban"},
|
||
{9, 3, 4, 1, "通知模板", "/notify_tpl/list", "fa-bullhorn"},
|
||
{10, 4, 1, 1, "执行日志", "/task_log/list", "fa-file"},
|
||
{11, 5, 1, 1, "权限因子", "/auth/index", "fa-key"},
|
||
{12, 5, 2, 1, "角色管理", "/role/list", "fa-group"},
|
||
{13, 5, 3, 1, "管理员管理", "/admin/list", "fa-user"},
|
||
}
|
||
for _, r := range rows {
|
||
g.DB().Exec(ctx,
|
||
"INSERT INTO pp_uc_auth(id, auth_name, auth_url, pid, sort, icon, is_show, status, create_time, update_time) VALUES(?,?,?,?,?,?,?,1,?,?)",
|
||
r.id, r.name, r.url, r.pid, r.sort, r.icon, r.isShow, now, now,
|
||
)
|
||
}
|
||
}
|
||
|
||
// seedBusinessData 种子数据:业务数据(任务分组 + 定时任务)
|
||
func seedBusinessData(ctx context.Context) {
|
||
now := time.Now().Unix()
|
||
|
||
// ---- 任务分组 ----
|
||
var groupCount int
|
||
_ = g.DB().GetScan(ctx, &groupCount, "SELECT COUNT(*) FROM pp_task_group")
|
||
if groupCount == 0 {
|
||
g.DB().Exec(ctx,
|
||
`INSERT INTO pp_task_group(id, group_name, description, create_id, update_id, create_time, update_time, status)
|
||
VALUES(?,?,?,?,?,?,?,?)`,
|
||
1, "数据引擎", "Data Engine 定时同步任务", 1, 1, now, now, 1,
|
||
)
|
||
g.Log().Info(ctx, "已创建默认任务分组: 数据引擎")
|
||
}
|
||
|
||
// ---- 定时任务 ----
|
||
var taskCount int
|
||
_ = g.DB().GetScan(ctx, &taskCount, "SELECT COUNT(*) FROM pp_task")
|
||
if taskCount == 0 {
|
||
seedTasks(ctx, now)
|
||
}
|
||
}
|
||
|
||
// seedTasks 种子定时任务
|
||
func seedTasks(ctx context.Context, now int64) {
|
||
type taskSeed struct {
|
||
id int
|
||
group_id int
|
||
server_ids, task_name, description, cron_spec string
|
||
concurrent int
|
||
task_type, command, url, method, headers, body string
|
||
timeout, status int
|
||
}
|
||
rows := []taskSeed{
|
||
{46, 1, "", "数据引擎-补偿扫描", "",
|
||
"0 */5 * * * *", 0,
|
||
"http", "", "http://host.docker.internal:3013/sync/ctrl/compensate", "POST",
|
||
`{"Content-Type":"application/json"}`, "",
|
||
600, 0},
|
||
{48, 1, "", "腾讯广告-账户列表(account_relation)", "",
|
||
"0 0 */6 * * *", 0,
|
||
"http", "", "http://host.docker.internal:3013/sync/ctrl/trigger", "POST",
|
||
`{"Content-Type":"application/json"}`,
|
||
`{"platformCode":"tencent","interfaceCode":"account_relation","fullSync":false}`,
|
||
1800, 0},
|
||
{49, 1, "", "腾讯广告-图片素材(image)", "",
|
||
"0 0 * * * *", 0,
|
||
"http", "", "http://host.docker.internal:3013/sync/ctrl/trigger", "POST",
|
||
`{"Content-Type":"application/json"}`,
|
||
`{"platformCode":"tencent","interfaceCode":"image","fullSync":false}`,
|
||
3600, 0},
|
||
{50, 1, "", "腾讯广告-视频素材(video)", "",
|
||
"0 0 * * * *", 0,
|
||
"http", "", "http://host.docker.internal:3013/sync/ctrl/trigger", "POST",
|
||
`{"Content-Type":"application/json"}`,
|
||
`{"platformCode":"tencent","interfaceCode":"video","fullSync":false}`,
|
||
3600, 0},
|
||
{51, 1, "", "腾讯广告-音频素材(audio)", "",
|
||
"0 0 */6 * * *", 0,
|
||
"http", "", "http://host.docker.internal:3013/sync/ctrl/trigger", "POST",
|
||
`{"Content-Type":"application/json"}`,
|
||
`{"platformCode":"tencent","interfaceCode":"audio","fullSync":false}`,
|
||
1800, 0},
|
||
{52, 1, "", "腾讯广告-Token刷新", "",
|
||
"0 0 3 * * *", 0,
|
||
"http", "", "http://host.docker.internal:3013/sync/ctrl/refreshToken", "POST",
|
||
`{"Content-Type":"application/json"}`,
|
||
`{"platformCode":"tencent"}`,
|
||
0, 1},
|
||
{53, 1, "", "CID - 图片批量送检", "自动扫描待校验图片并提交到易盾检测",
|
||
"0/30 * * * * *", 0,
|
||
"http", "", "http://host.docker.internal:3001/material/verify/controller/batch-verify-image", "POST",
|
||
`{"Content-Type":"application/json"}`,
|
||
"{limit:10}",
|
||
120, 0},
|
||
{54, 1, "", "CID - 视频批量送检", "自动扫描待校验视频并提交到易盾检测",
|
||
"15/30 * * * * *", 0,
|
||
"http", "", "http://host.docker.internal:3001/material/verify/controller/batch-verify-video", "POST",
|
||
`{"Content-Type":"application/json"}`,
|
||
"{limit:10}",
|
||
120, 0},
|
||
{55, 1, "", "CID - 检测结果轮询", "自动查询易盾检测结果并更新到素材表",
|
||
"0 * * * * *", 0,
|
||
"http", "", "http://host.docker.internal:3001/yidun/callback/controller/poll-all-results", "POST",
|
||
"{}", "",
|
||
120, 1},
|
||
}
|
||
|
||
for _, r := range rows {
|
||
_, err := g.DB().Exec(ctx,
|
||
`INSERT INTO pp_task(
|
||
id, group_id, server_ids, server_type, task_name, description, cron_spec,
|
||
concurrent, task_type, command, url, method, headers, body,
|
||
timeout, execute_times, prev_time, status,
|
||
is_notify, notify_type, notify_tpl_id, notify_user_ids,
|
||
create_id, update_id, create_time, update_time
|
||
) VALUES(?,?,?,0,?,?,?,?,?,?,?,?,?,?,?,0,0,?, 0,0,0,'', 1,1,?,?)`,
|
||
r.id, r.group_id, r.server_ids, r.task_name, r.description, r.cron_spec,
|
||
r.concurrent, r.task_type, r.command, r.url, r.method, r.headers, r.body,
|
||
r.timeout, r.status,
|
||
now, now,
|
||
)
|
||
if err != nil {
|
||
g.Log().Warningf(ctx, "种子任务插入失败 (id=%d): %v", r.id, err)
|
||
}
|
||
}
|
||
g.Log().Infof(ctx, "已创建 %d 个种子定时任务", len(rows))
|
||
}
|