Files

302 lines
13 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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))
}