gcli2api / internal /panel /taskcenter.go
a3216's picture
chore: 同步到上游 1.12.0-panel + 凭证同步/独立启动器/保活
6d60378 verified
Raw History Blame Contribute Delete
17.1 kB
// taskcenter.go 面板「任务中心」:全账号任务扫描 + 执行队列(可配并发)+
// 成长任务队列。
//
// 语义:
// - 扫描(scan_all):并发拉取每账号的成长任务列表(默认+小程序口径),
// 汇总出"未完成且可自动化"的待办清单(只读,不执行)。
// - 执行队列(run_queue + queue):把待办项按账号分组排队执行——账号内
// 串行(复用 per-account 锁,与单任务/一键完成互斥),账号间并发
// (concurrency 信号量限制,默认 1)。队列状态可轮询。
package panel
import (
"encoding/json"
"fmt"
"log"
"net/http"
"sort"
"sync"
"time"
"github.com/linguo2625469/workbuddy2api-panel/internal/auth"
"github.com/linguo2625469/workbuddy2api-panel/internal/upstream"
)
// ---------------------------------------------------------------------------
// 扫描(只读)
// ---------------------------------------------------------------------------
// scanAccountItem 单账号扫描结果。
type scanAccountItem struct {
UID string `json:"uid"`
Nickname string `json:"nickname"`
Growth []upstream.Task `json:"growth,omitempty"`
GrowthErr string `json:"growth_error,omitempty"`
}
// growthPending 任务是否"未完成且可自动化"。
func growthPending(t upstream.Task) bool {
if t.Claimed {
return false
}
// 上游锁定的任务不出待办:Sequential 族每日零点解锁一环,刚做完上一环时
// 下一环以下发但 locked 形态出现在列表里——扫进队列只会 accept 不落账报
// 失败(每日锁定窗口),零点解锁后自然回到待办。其余 locked(上游未开放)
// 同语义:不该被自动化尝试。
if t.Locked {
return false
}
if t.Target > 0 && t.Current >= t.Target {
// 达标未领:也入队(队列执行后会自动领)——但仅限有自动化动作的任务,
// 否则队列执行时会因 autoActionFor 为 nil 直接报错。
return autoActionFor(t.TaskCode) != nil
}
return autoActionFor(t.TaskCode) != nil
}
// tasksScanAll 扫描全部账号:成长任务(未完成+可自动化,含 mp 口径合并)。
// 只读操作,并发拉取(账号数个位数)。
func (p *Panel) tasksScanAll(w http.ResponseWriter, r *http.Request) {
states := p.cfg.Pool.List()
items := make([]scanAccountItem, len(states))
var wg sync.WaitGroup
for i, st := range states {
if st.Disabled {
continue
}
wg.Add(1)
go func(i int, uid string) {
defer wg.Done()
a := p.cfg.Pool.AuthByUID(uid)
if a == nil {
return
}
it := &items[i]
it.UID, it.Nickname = uid, a.Nickname
// D4 门控:global 账号无 CN 成长任务体系,不发起任何上游调用。
if a.IsGlobal() {
return
}
if tasks, err := p.cfg.Upstream.ListTasks(a); err != nil {
it.GrowthErr = err.Error()
} else {
for _, t := range tasks {
if growthPending(t) {
it.Growth = append(it.Growth, t)
}
}
}
// 小程序口径任务(school_season 校园日 / Sequential_Tasks_1 小程序首对话)
// 仅在 mp 头列表下发,与默认口径不重叠——合并进待办列表;mp 列表失败
// 静默(无 mp 任务的部署/活动结束时零影响)。
if mpTasks, err := p.cfg.Upstream.ListTasksMP(a); err == nil {
seen := map[string]bool{}
for _, t := range it.Growth {
seen[t.TaskCode] = true
}
for _, t := range mpTasks {
if growthPending(t) && !seen[t.TaskCode] {
it.Growth = append(it.Growth, t)
}
}
}
}(i, st.UID)
}
wg.Wait()
pending := 0
for _, it := range items {
pending += len(it.Growth)
}
log.Printf("panel: 队列扫描完成:全部账号待办 %d 项", pending)
writeJSON(w, http.StatusOK, map[string]any{"ok": true, "accounts": items, "pending_count": pending})
}
// ---------------------------------------------------------------------------
// 执行队列
// ---------------------------------------------------------------------------
// queueItem 队列执行单元。
type queueItem struct {
UID string `json:"uid"`
Nickname string `json:"nickname"`
Kind string `json:"kind"` // growth
Code string `json:"code"`
Status string `json:"status"` // pending | running | done | skipped | error
Message string `json:"message,omitempty"`
}
// queueState 队列运行状态。Seq 每次启动 +1——前端只渲染"自己启动的那一轮",
// 执行结束后的残留 items 不会覆盖后续的扫描结果视图。
type queueState struct {
mu sync.Mutex
running bool
startedAt time.Time
items []queueItem
conc int
seq int
}
// Panel 队列字段在 Panel 结构体上(panel.go)由 initQueue 惰性初始化;
// 这里集中访问器,避免改动 New 构造链。
func (p *Panel) queue() *queueState {
p.queueOnce.Do(func() { p.q = &queueState{} })
return p.q
}
// tasksRunQueue 启动执行队列:{concurrency:1-4, growth:bool, school:bool}。
// 先做一次扫描,把全部待办项排队(growth 按账号内 autoActions 顺序执行,
// school 逐账号跑闭环),账号内串行、账号间受并发信号量约束。
func (p *Panel) tasksRunQueue(w http.ResponseWriter, r *http.Request) {
var body struct {
Concurrency int `json:"concurrency"`
Growth bool `json:"growth"`
}
_ = json.NewDecoder(r.Body).Decode(&body)
if !body.Growth {
body.Growth = true
}
if body.Concurrency < 1 {
body.Concurrency = 1
}
if body.Concurrency > 4 {
body.Concurrency = 4
}
started, total, seq, msg := p.startGrowthQueue(body.Concurrency, body.Growth)
switch {
case seq == -1:
writeErr(w, http.StatusConflict, msg)
case !started:
writeJSON(w, http.StatusOK, map[string]any{"ok": true, "started": false, "message": msg})
default:
writeJSON(w, http.StatusOK, map[string]any{"ok": true, "started": true, "total": total, "seq": seq})
}
}
// startGrowthQueue 扫描全部账号待办并启动队列(HTTP「执行全部待办」与调度器
// growth 时点共用核心)。返回 (started, total, seq, msg):seq==-1 表示队列
// 已在执行(冲突);started=false 时 msg 为无可执行待办的说明。并发夹取
// [1,4];growth 开关同 HTTP 入参语义。
func (p *Panel) startGrowthQueue(concurrency int, growth bool) (started bool, total int, seq int, msg string) {
if concurrency < 1 {
concurrency = 1
}
if concurrency > 4 {
concurrency = 4
}
q := p.queue()
q.mu.Lock()
if q.running {
q.mu.Unlock()
return false, 0, -1, "队列正在执行中(可在任务中心查看进度)"
}
// 先占位:扫描(数秒级网络耗时)期间若并发再次触发,直接命中上面的 running
// 判拒,避免两个 goroutine 同时启动互相覆盖 q.items/q.seq。无待办时回滚。
q.running = true
q.startedAt = time.Now()
q.mu.Unlock()
// 扫描待办(复用扫描逻辑的拉取部分)。
states := p.cfg.Pool.List()
var accts []queueAccount
var wg sync.WaitGroup
var mu sync.Mutex
for _, st := range states {
if st.Disabled {
continue
}
a := p.cfg.Pool.AuthByUID(st.UID)
if a == nil {
continue
}
wg.Add(1)
go func(a *auth.Auth) {
defer wg.Done()
one := queueAccount{a: a}
// D4 门控:global 账号无 CN 成长任务体系,不发起任何上游调用。
if a.IsGlobal() {
return
}
if growth {
if tasks, err := p.cfg.Upstream.ListTasks(a); err == nil {
for _, t := range tasks {
if growthPending(t) {
one.grow = append(one.grow, t)
}
}
// 合并小程序口径待办(与 tasksScanAll 同口径:mp 列表是默认口径
// 超集,按 code 去重;失败静默)。此前此处漏合并——扫描显示
// mp 待办而队列报"无可执行待办"。
if mpTasks, mpErr := p.cfg.Upstream.ListTasksMP(a); mpErr == nil {
seen := map[string]bool{}
for _, t := range one.grow {
seen[t.TaskCode] = true
}
for _, t := range mpTasks {
if growthPending(t) && !seen[t.TaskCode] {
one.grow = append(one.grow, t)
}
}
}
sort.Slice(one.grow, func(i, j int) bool { // 按 autoActions 顺序(依赖前置)
return autoActionIndex(one.grow[i].TaskCode) < autoActionIndex(one.grow[j].TaskCode)
})
}
}
if len(one.grow) > 0 {
mu.Lock()
accts = append(accts, one)
mu.Unlock()
}
}(a)
}
wg.Wait()
// 组装队列(账号分组,保持顺序)。
var items []queueItem
for _, one := range accts {
for _, t := range one.grow {
items = append(items, queueItem{UID: one.a.UID, Nickname: one.a.Nickname, Kind: "growth", Code: t.TaskCode, Status: "pending"})
}
}
if len(items) == 0 {
log.Printf("panel: 队列启动:无可执行待办(全部账号任务已完成)")
q.mu.Lock()
q.running = false
q.startedAt = time.Time{}
q.mu.Unlock()
return false, 0, 0, "全部账号没有待办任务"
}
q.mu.Lock()
q.items = items
q.conc = concurrency
q.seq++
seq = q.seq
q.mu.Unlock()
go p.runQueueItems(accts, items, concurrency)
log.Printf("panel: 队列启动:%d 项(并发 %d,成长 %v)", len(items), concurrency, growth)
return true, len(items), seq, ""
}
// RunGrowthQueueOnce 调度器 growth 时点回调(sch.SetGrowthHook 挂载):与
// 「执行全部待办」按钮完全同管线(成长,串行并发 1)。Sequential 族
// 每日零点解锁一环,此前只能手动扫描推进;此回调让链条每天自动走一环。
// 异步执行(startGrowthQueue 启动 goroutine 即返),已在跑/无待办安全跳过。
func (p *Panel) RunGrowthQueueOnce() {
started, total, _, _ := p.startGrowthQueue(1, true)
if started {
log.Printf("panel: 定时成长任务队列已启动(%d 项)", total)
}
}
// runQueueItems 队列执行主体:按账号分组,账号内串行(per-account 锁),
// 账号间并发(信号量)。每项结果写回队列状态。
func (p *Panel) runQueueItems(accts []queueAccount, items []queueItem, concurrency int) {
q := p.queue()
defer func() {
q.mu.Lock()
q.running = false
q.mu.Unlock()
log.Printf("panel: 队列执行结束(共 %d 项)", len(items))
}()
sem := make(chan struct{}, concurrency)
var wg sync.WaitGroup
for _, one := range accts {
wg.Add(1)
go func(one queueAccount) {
defer wg.Done()
sem <- struct{}{}
defer func() { <-sem }()
// per-account 互斥:与单任务/一键完成共用一把锁。
if !p.tryLockAccount(one.a.UID) {
p.queueSet(q, one.a.UID, func(it *queueItem) {
it.Status, it.Message = "skipped", "该账号有其它任务动作在执行,跳过"
})
return
}
defer p.unlockAccount(one.a.UID)
// 前置:批量接受尚未接受的任务。上游对 not_accepted 的任务不计数——
// 面板「一键完成」一直有这步,队列路径此前漏了(表现为上报 200 但进度
// 一直 not_accepted、无法领奖)。失败不阻塞(行为事件才是进度判据)。
if accepted := p.acceptPendingTasks(one.a); accepted > 0 {
time.Sleep(reportGap) // 给上游状态流转留时间
}
for i := range q.items {
uid, kind, code := q.snapshotAt(i)
if uid != one.a.UID {
continue
}
p.queueMarkAt(i, "running", "")
var msg string
var err error
switch kind {
case "growth":
msg, err = p.runGrowthQueued(one.a, code)
}
if err != nil {
p.queueMarkAt(i, "error", err.Error())
} else {
p.queueMarkAt(i, "done", msg)
}
time.Sleep(reportGap) // 项间节流
}
}(one)
}
wg.Wait()
}
// queueAccount 队列执行的账号单元(runQueueItems 参数)。
type queueAccount struct {
a *auth.Auth
grow []upstream.Task
}
// snapshotAt 锁内读条目三元组(避免锁外持有指针)。
func (q *queueState) snapshotAt(i int) (uid, kind, code string) {
q.mu.Lock()
defer q.mu.Unlock()
return q.items[i].UID, q.items[i].Kind, q.items[i].Code
}
// queueMarkAt 按索引更新队列条目状态(条目数组固定不再增删)。
func (p *Panel) queueMarkAt(i int, status, msg string) {
q := p.queue()
q.mu.Lock()
q.items[i].Status, q.items[i].Message = status, msg
q.mu.Unlock()
}
// queueSet 按 uid 批量改状态。
func (p *Panel) queueSet(q *queueState, uid string, fn func(*queueItem)) {
q.mu.Lock()
defer q.mu.Unlock()
for i := range q.items {
if q.items[i].UID == uid {
fn(&q.items[i])
}
}
}
// acceptPendingTasks 批量接受该账号未接受的任务,返回接受的个数(失败返回 0 不阻塞)。
func (p *Panel) acceptPendingTasks(a *auth.Auth) int {
tasks, err := p.cfg.Upstream.ListTasks(a)
if err != nil {
return 0
}
var codes []string
for _, t := range tasks {
if !t.Claimed && !t.Locked && t.AcceptStatus != "accepted" && t.AcceptStatus != "completed" {
codes = append(codes, t.TaskCode)
}
}
if len(codes) == 0 {
return 0
}
if err := p.cfg.Upstream.AcceptTasks(a, codes); err != nil {
log.Printf("panel: 队列 accept uid=%s: %v(不阻塞)", a.UID, err)
return 0
}
log.Printf("panel: 队列 accept uid=%s: 已接受 %d 个任务", a.UID, len(codes))
return len(codes)
}
// runGrowthQueued 执行单个成长任务(动作 + 回读 + 自动领奖;与
// accountTaskAuto 同语义,结果以文字返回)。
func (p *Panel) runGrowthQueued(a *auth.Auth, code string) (string, error) {
act := autoActionFor(code)
if act == nil {
return "", fmt.Errorf("任务 %s 无自动动作", code)
}
// taskByCode 已双口径(mp 专属码自动回落 mp 列表)。
before, err := p.taskByCode(a, code)
if err != nil {
return "", err
}
if before == nil {
return "该账号无此任务", nil
}
isMP := isMPTaskCode(code)
if before.Claimed {
return "已完成(已领取)", nil
}
msg, err := act.run(p, a)
if err != nil {
return "", err
}
var after *upstream.Task
if isMP {
after, _ = p.taskByCodeMP(a, code)
} else {
after, _ = p.taskByCodeWaiting(a, code)
}
if after != nil && after.Claimable {
var credit, energy int64
var cerr error
if isMP {
credit, energy, cerr = p.cfg.Upstream.ClaimRewardMP(a, code)
} else {
credit, energy, cerr = p.cfg.Upstream.ClaimReward(a, code)
}
if cerr == nil && (credit > 0 || energy > 0) {
msg += fmt.Sprintf(";自动领奖 +%d 分 +%d 能", credit, energy)
}
}
if after != nil {
msg += "(进度 " + taskProgressText(after) + ")"
}
log.Printf("panel: 队列 growth uid=%s code=%s: %s", a.UID, code, msg)
return msg, nil
}
// tasksQueueStatus 队列状态(轮询用)。
func (p *Panel) tasksQueueStatus(w http.ResponseWriter, r *http.Request) {
q := p.queue()
q.mu.Lock()
defer q.mu.Unlock()
items := make([]queueItem, len(q.items))
copy(items, q.items)
writeJSON(w, http.StatusOK, map[string]any{
"running": q.running,
"total": len(items),
"conc": q.conc,
"started": !q.startedAt.IsZero(),
"started_at": q.startedAt,
"seq": q.seq,
"items": items,
})
}
// schoolVouchers 我的券码:逐 CN 账号查开学季 /vouchers(3 并发,与 packages
// 同款限流),失败只在对应账号标 error。global 账号无开学季,不发上游调用。
func (p *Panel) schoolVouchers(w http.ResponseWriter, r *http.Request) {
accts := p.cfg.Pool.List()
type row struct {
UID string `json:"uid"`
Nickname string `json:"nickname"`
Vouchers []upstream.SchoolVoucher `json:"vouchers"`
Err string `json:"error,omitempty"`
}
out := make([]row, len(accts))
sem := make(chan struct{}, 3)
var wg sync.WaitGroup
for i, st := range accts {
if st.Disabled {
continue // 未占位,行末统一压掉
}
a := p.cfg.Pool.AuthByUID(st.UID)
if a == nil {
continue
}
wg.Add(1)
go func(i int, a *auth.Auth) {
defer wg.Done()
sem <- struct{}{}
defer func() { <-sem }()
it := row{UID: a.UID, Nickname: a.Nickname}
switch {
case a.IsGlobal():
it.Err = "global realm(无开学季活动)"
default:
vs, err := p.cfg.Upstream.SchoolVouchers(a)
if err != nil {
it.Err = err.Error()
} else {
it.Vouchers = vs
}
}
out[i] = it
}(i, a)
}
wg.Wait()
res := make([]row, 0, len(out))
for _, it := range out {
if it.UID != "" {
res = append(res, it)
}
}
writeJSON(w, http.StatusOK, map[string]any{"accounts": res})
}