onesvm-browser-server/server/internal/scheduler/loops.go
chii eb972dfa93 feat: 落地 browser-server 控制面并打通 mgr1 海外订阅
单二进制三角色 + Dock 适配器 + Swarm stack 达到可部署态;mgr1 实测订阅经 central-proxy bootstrap,探活 alive=41/52。

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-09-02 15:05:12 +08:00

190 lines
5.2 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.

// loops.go:worker 池 / reaper / proxy 探活循环(Core 生命周期部分)。
package scheduler
import (
"context"
"time"
"onesvm.com/onesvm/browser-server/internal/contract"
"onesvm.com/onesvm/browser-server/internal/store"
)
// workerLoop worker 主循环:快通道唤醒 + 兜底轮询 ClaimNext(原子抢单)。
func (c *Core) workerLoop(ctx context.Context, idx int) {
defer c.wg.Done()
workerName := "worker-" + itoa(idx)
poll := time.NewTicker(500 * time.Millisecond) // 兜底轮询(chan 唤醒失败/回队任务)
defer poll.Stop()
for {
select {
case <-ctx.Done():
return
case <-c.stopCh:
return
case <-poll.C:
}
c.claimAndRun(ctx, workerName)
}
}
// claimAndRun 抢一单执行一单。
func (c *Core) claimAndRun(ctx context.Context, workerName string) {
job, err := c.db.ClaimNext(workerName, LeaseFor)
if err != nil {
if err != store.ErrNotFound {
c.log.Printf("ClaimNext 错误: %v", err)
}
return
}
c.st.RunningG.Add(1)
defer c.st.RunningG.Add(-1)
c.runJob(ctx, job)
}
// reaperLoop 每 5s 扫租约过期(A3.1:lease_until < now AND status=running →
// attempts+1(ClaimNext 已计)→ attempts≤2 且瞬时错误重入队(指数退避)→
// 否则死信;store.ReapExpired 已实现该语义)。
func (c *Core) reaperLoop(ctx context.Context) {
defer c.wg.Done()
t := time.NewTicker(ReaperEvery)
defer t.Stop()
for {
select {
case <-ctx.Done():
return
case <-c.stopCh:
return
case <-t.C:
n, err := c.db.ReapExpired(MaxAttempts)
if err != nil {
c.log.Printf("reaper 错误: %v", err)
continue
}
if n > 0 {
c.log.Printf("reaper 收割 %d 个租约过期任务(退避回队/死信)", n)
}
}
}
}
// healthWatchLoop 适配器健康重探(ITER-1 F2 + ITER-2 冷启动收敛,fail-w5-smoke-iter1):
// 适配器 Init 为一次性探活,与引擎容器启动存在竞态(DNS/端口未就绪 → healthy=false
// 永不恢复 → 任务 upstream→dead,T6 瞬态失败永久化)。本循环即「摘除→恢复回表」的
// 完整闭环(design-arch §4.5 条款 2 + §3.5 降级条款)。并入 Core.Start 生命周期。
//
// ITER-2 收敛窗口:冷启动阶段(healthWatchBootProbes 次内)按 1s 短间隔只重探
// Health().OK==false 的适配器(健康者跳过,省流量;Init 幂等无害);全部健康后
// 切入 healthWatchEvery 常规周期。冷启动窗口 10×1s + 探活 5s 超时,覆盖 compose
// 同时 up 的 python 容器 1–2s 冷启;生产引擎启动亦不会超此窗口。10 次仍不健康
// 则退 30s 常规周期(引擎故障属常态摘除,不该空转打探活)。
func (c *Core) healthWatchLoop(ctx context.Context) {
defer c.wg.Done()
// 冷启动阶段:最多 10 次、1s 间隔,只重探不健康适配器。
for i := 0; i < healthWatchBootProbes; i++ {
// 启动即重探一次(对齐 proxyWatchLoop「启动即探一次」先例)。
_ = c.reg.InitAll(ctx) // 幂等;全量首探收敛首轮健康位
if c.allAdaptersHealthy() {
break
}
select {
case <-ctx.Done():
return
case <-c.stopCh:
return
case <-time.After(healthWatchBootInterval):
c.reprobeUnhealthy(ctx)
}
}
// 全部健康(或 10 次未收敛)→ 切 30s 常规周期(不健康者周期重探,健康者 Init 幂等无害)。
t := time.NewTicker(healthWatchEvery)
defer t.Stop()
for {
select {
case <-ctx.Done():
return
case <-c.stopCh:
return
case <-t.C:
initErrs := c.reg.InitAll(ctx)
for name, e := range initErrs {
c.log.Printf("适配器 %s 健康重探错误(保持摘除): %v", name, e)
}
}
}
}
// allAdaptersHealthy 全部适配器健康判定(冷启动收敛用)。
func (c *Core) allAdaptersHealthy() bool {
for _, a := range c.reg.All() {
if !a.Health().OK {
return false
}
}
return true
}
// reprobeUnhealthy 只对不健康适配器重跑 Init(冷启动阶段省流量;Init 幂等)。
func (c *Core) reprobeUnhealthy(ctx context.Context) {
for _, a := range c.reg.All() {
if a.Health().OK {
continue
}
_ = a.Init(ctx)
}
}
// proxyWatchLoop proxymanager 探活:可达位维护(overseas fail-closed 依据,O3)。
func (c *Core) proxyWatchLoop(ctx context.Context) {
defer c.wg.Done()
t := time.NewTicker(proxyProbeInterval)
defer t.Stop()
// 启动即探一次。
c.proxyOK.Store(c.proxy.Ping(ctx))
for {
select {
case <-ctx.Done():
return
case <-c.stopCh:
return
case <-t.C:
ok := c.proxy.Ping(ctx)
was := c.proxyOK.Swap(ok)
if ok != was {
if ok {
c.log.Printf("proxymanager 恢复可达(overseas 通道恢复)")
} else {
c.log.Printf("proxymanager 不可达(overseas 任务将 fail-closed 排队等待)")
}
}
}
}
}
// itoa 简单 int → string(避免 fmt 进热路径)。
func itoa(n int) string {
if n == 0 {
return "0"
}
neg := n < 0
if neg {
n = -n
}
var buf [20]byte
i := len(buf)
for n > 0 {
i--
buf[i] = byte('0' + n%10)
n /= 10
}
if neg {
i--
buf[i] = '-'
}
return string(buf[i:])
}
// retryableJob 判断任务错误是否瞬时(重试入口断言)。
func retryableJob(code string) bool { return contract.Retryable(code) }
var _ = store.ErrNotFound