supervisor.go 9.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391
  1. package main
  2. import (
  3. "bufio"
  4. "errors"
  5. "fmt"
  6. "io"
  7. "os/exec"
  8. "strconv"
  9. "strings"
  10. "sync"
  11. "time"
  12. )
  13. // supervisor.go —— dsh web 守护核心:端口预检、隐藏启动、指数退避重启、状态快照。
  14. // healthyUptime 稳定运行判定阈值:超过它则清零失败计数(退避复位)。
  15. const healthyUptime = 30 * time.Second
  16. // Status 服务状态快照。
  17. type Status struct {
  18. Running bool
  19. Ready bool // 已捕获带 token 的访问地址,Web 界面可打开
  20. PID int
  21. Restarts int
  22. Uptime time.Duration
  23. LastError string
  24. URL string
  25. Updating bool
  26. External bool // 端口被外部实例占用(非本程序启动)
  27. }
  28. // Supervisor 守护 dsh web。
  29. type Supervisor struct {
  30. mu sync.Mutex
  31. probeMu sync.Mutex // 串行化入口探测(见 probe.go)
  32. cfg *Config
  33. log *Logger
  34. pid int
  35. running bool
  36. desired bool
  37. restarts int
  38. failures int
  39. startedAt time.Time
  40. lastError string
  41. updating bool
  42. external bool
  43. binCache string
  44. nodeCache string
  45. webURL string // dsh web 输出的完整 URL(含访问 token,仅存内存)
  46. autoMu sync.Mutex
  47. autoLoop bool // 定期更新循环是否已启动(进程生命周期内单例)
  48. }
  49. func NewSupervisor(cfg *Config, log *Logger) *Supervisor {
  50. return &Supervisor{cfg: cfg, log: log}
  51. }
  52. // Start 设置期望运行并拉起守护循环。
  53. func (s *Supervisor) Start() {
  54. s.mu.Lock()
  55. if s.desired {
  56. s.mu.Unlock()
  57. return
  58. }
  59. s.desired = true
  60. s.failures = 0
  61. s.mu.Unlock()
  62. go s.monitor()
  63. }
  64. // Stop 同步停止:结束进程树并等待其真正退出(最多 8 秒)。
  65. func (s *Supervisor) Stop() {
  66. s.mu.Lock()
  67. s.desired = false
  68. pid := s.pid
  69. s.webURL = ""
  70. s.mu.Unlock()
  71. if pid > 0 {
  72. s.log.Printf("停止 dsh web pid=%d", pid)
  73. _ = killTree(pid)
  74. s.waitStopped(pid, 8*time.Second)
  75. }
  76. }
  77. // SetAutoUpdate 设置定期更新开关(加锁 + 立即持久化;开启时确保循环已启动)。
  78. func (s *Supervisor) SetAutoUpdate(enabled bool) error {
  79. s.mu.Lock()
  80. s.cfg.AutoUpdate = enabled
  81. snapshot := *s.cfg
  82. s.mu.Unlock()
  83. if enabled {
  84. s.startAutoUpdate()
  85. }
  86. return saveConfig(&snapshot)
  87. }
  88. // SetAutoStart 设置开机自启开关(加锁 + 立即持久化;注册表由调用方处理)。
  89. func (s *Supervisor) SetAutoStart(enabled bool) error {
  90. s.mu.Lock()
  91. s.cfg.AutoStart = enabled
  92. snapshot := *s.cfg
  93. s.mu.Unlock()
  94. return saveConfig(&snapshot)
  95. }
  96. // UpdateAndRestart 先停止服务再更新,避免 Windows 上 native 模块文件被占用导致安装失败;
  97. // 无论更新成功与否都恢复到原有运行意图。已是最新版本时既不安装也不重启。
  98. func (s *Supervisor) UpdateAndRestart() (string, error) {
  99. if fresh, current := s.upToDate(); fresh {
  100. s.log.Printf("更新:已是最新(%s),不重启服务", current)
  101. return "已是最新版本 " + current, nil
  102. }
  103. wasDesired := s.isDesired()
  104. if wasDesired {
  105. s.Stop()
  106. }
  107. out, err := s.Update()
  108. if wasDesired {
  109. s.Start()
  110. }
  111. return out, err
  112. }
  113. // Restart 重启并等待端口释放,避免新旧实例争抢端口。
  114. func (s *Supervisor) Restart() {
  115. s.log.Printf("重启 dsh web")
  116. s.Stop()
  117. s.waitPortFree(5 * time.Second)
  118. s.Start()
  119. }
  120. func (s *Supervisor) Status() Status {
  121. s.mu.Lock()
  122. defer s.mu.Unlock()
  123. url := s.baseURL()
  124. if s.webURL != "" {
  125. url = s.webURL // 优先使用带 token 的地址
  126. }
  127. var uptime time.Duration
  128. if s.running && !s.startedAt.IsZero() {
  129. uptime = time.Since(s.startedAt)
  130. }
  131. return Status{
  132. Running: s.running,
  133. Ready: s.running && s.webURL != "",
  134. PID: s.pid,
  135. Restarts: s.restarts,
  136. Uptime: uptime,
  137. LastError: s.lastError,
  138. URL: url,
  139. Updating: s.updating,
  140. External: s.external,
  141. }
  142. }
  143. func (s *Supervisor) baseURL() string {
  144. host := s.cfg.WebHost
  145. if host == "" {
  146. host = "127.0.0.1"
  147. }
  148. return fmt.Sprintf("http://%s:%d", host, s.cfg.WebPort)
  149. }
  150. func (s *Supervisor) setError(err error) {
  151. s.mu.Lock()
  152. s.lastError = err.Error()
  153. s.mu.Unlock()
  154. s.log.Printf("守护错误: %v", err)
  155. }
  156. // maxBackoff 重启退避上限。
  157. const maxBackoff = 2 * time.Minute
  158. // backoff 指数退避:base * 2^failures,封顶 maxBackoff。
  159. // 循环内提前 break,既保证封顶语义,也避免超大 failures 造成数值溢出。
  160. func (s *Supervisor) backoff() time.Duration {
  161. s.mu.Lock()
  162. failures := s.failures
  163. base := s.cfg.RestartDelaySec
  164. s.mu.Unlock()
  165. if base <= 0 {
  166. base = 5
  167. }
  168. d := time.Duration(base) * time.Second
  169. for i := 0; i < failures; i++ {
  170. if d >= maxBackoff {
  171. break
  172. }
  173. d *= 2
  174. }
  175. if d > maxBackoff {
  176. d = maxBackoff
  177. }
  178. return d
  179. }
  180. // portBusyBackoff 端口占用时的温和退避(上限 30 秒):不是故障,只是等待外部实例释放。
  181. func (s *Supervisor) portBusyBackoff() time.Duration {
  182. d := s.backoff()
  183. if d > 30*time.Second {
  184. d = 30 * time.Second
  185. }
  186. return d
  187. }
  188. // monitor 守护循环:端口预检 -> 启动 -> 等待退出 -> 期望运行时退避重启。
  189. func (s *Supervisor) monitor() {
  190. for {
  191. if !s.isDesired() {
  192. return
  193. }
  194. // 端口预检:已有实例监听时不再启动,避免端口冲突导致的重启风暴
  195. if isPortOpen(s.cfg.WebHost, s.cfg.WebPort) {
  196. s.mu.Lock()
  197. s.running = false
  198. s.pid = 0
  199. s.external = true
  200. s.lastError = "端口已被占用(可能已有 dsh web 在运行)"
  201. s.mu.Unlock()
  202. s.log.Printf("端口 %s 已被占用,跳过启动(等待释放)", s.baseURL())
  203. s.bumpFailure()
  204. s.sleepInterruptible(s.portBusyBackoff())
  205. continue
  206. }
  207. s.mu.Lock()
  208. s.external = false
  209. s.mu.Unlock()
  210. exe, args, err := s.resolveCommand()
  211. if err != nil {
  212. s.setError(err)
  213. s.bumpFailure()
  214. s.sleepInterruptible(s.backoff())
  215. continue
  216. }
  217. cmd := exec.Command(exe, args...)
  218. hideWindow(cmd)
  219. stdout, errOut := cmd.StdoutPipe()
  220. stderr, errErr := cmd.StderrPipe()
  221. if errOut != nil || errErr != nil {
  222. s.setError(errors.New("无法创建输出管道"))
  223. s.bumpFailure()
  224. s.sleepInterruptible(s.backoff())
  225. continue
  226. }
  227. if err := cmd.Start(); err != nil {
  228. s.setError(err)
  229. s.bumpFailure()
  230. s.sleepInterruptible(s.backoff())
  231. continue
  232. }
  233. pid := cmd.Process.Pid
  234. startedAt := time.Now()
  235. s.mu.Lock()
  236. s.pid = pid
  237. s.running = true
  238. s.startedAt = startedAt
  239. s.lastError = ""
  240. s.webURL = ""
  241. desiredNow := s.desired
  242. s.mu.Unlock()
  243. s.log.Printf("dsh web 已启动 pid=%d: %s %s", pid, exe, strings.Join(args, " "))
  244. if !desiredNow {
  245. // 竞态窗口:停止请求发生在 cmd.Start 之后、pid 赋值之前。
  246. // 此时 Stop 读到的 pid 为 0 不会终止进程,这里补一次终止,避免残留。
  247. s.log.Printf("启动期间收到停止请求,立即结束 pid=%d", pid)
  248. _ = killTree(pid)
  249. }
  250. go s.pipeLog("out", stdout)
  251. go s.pipeLog("err", stderr)
  252. waitErr := cmd.Wait()
  253. ranFor := time.Since(startedAt)
  254. s.mu.Lock()
  255. s.running = false
  256. s.pid = 0
  257. stillDesired := s.desired
  258. if ranFor >= healthyUptime {
  259. s.failures = 0 // 稳定运行过,退避复位
  260. } else {
  261. s.failures++
  262. }
  263. if stillDesired {
  264. s.restarts++
  265. }
  266. s.mu.Unlock()
  267. s.log.Printf("dsh web 退出(err=%v,运行 %s),期望运行=%v", waitErr, ranFor.Round(time.Second), stillDesired)
  268. if !stillDesired {
  269. return
  270. }
  271. delay := s.backoff()
  272. s.log.Printf("%s 后自动重启", delay)
  273. s.sleepInterruptible(delay)
  274. }
  275. }
  276. func (s *Supervisor) isDesired() bool {
  277. s.mu.Lock()
  278. defer s.mu.Unlock()
  279. return s.desired
  280. }
  281. // bumpFailure 增加失败计数(触发退避放大)。
  282. func (s *Supervisor) bumpFailure() {
  283. s.mu.Lock()
  284. s.failures++
  285. s.mu.Unlock()
  286. }
  287. func (s *Supervisor) sleepInterruptible(d time.Duration) {
  288. deadline := time.Now().Add(d)
  289. for time.Now().Before(deadline) {
  290. if !s.isDesired() {
  291. return
  292. }
  293. time.Sleep(200 * time.Millisecond)
  294. }
  295. }
  296. func (s *Supervisor) waitStopped(pid int, timeout time.Duration) {
  297. deadline := time.Now().Add(timeout)
  298. for time.Now().Before(deadline) {
  299. if !processAlive(pid) {
  300. return
  301. }
  302. time.Sleep(100 * time.Millisecond)
  303. }
  304. s.log.Printf("等待 pid=%d 退出超时(继续)", pid)
  305. }
  306. func (s *Supervisor) waitPortFree(timeout time.Duration) {
  307. deadline := time.Now().Add(timeout)
  308. for time.Now().Before(deadline) {
  309. if !isPortOpen(s.cfg.WebHost, s.cfg.WebPort) {
  310. return
  311. }
  312. time.Sleep(150 * time.Millisecond)
  313. }
  314. }
  315. // pipeLog 把子进程输出写入日志(token 由 Logger 统一脱敏),并捕获访问地址。
  316. func (s *Supervisor) pipeLog(tag string, r io.ReadCloser) {
  317. if r == nil {
  318. return
  319. }
  320. sc := bufio.NewScanner(r)
  321. sc.Buffer(make([]byte, 0, 64*1024), 1024*1024)
  322. for sc.Scan() {
  323. line := sc.Text()
  324. s.log.Printf("[dsh %s] %s", tag, line)
  325. if url := extractHTTPURL(line); url != "" {
  326. s.mu.Lock()
  327. if url != s.webURL {
  328. s.webURL = url
  329. s.log.Printf("捕获 Web 访问地址(含 token,仅存内存);启动耗时 %s", time.Since(s.startedAt).Round(time.Second))
  330. }
  331. s.mu.Unlock()
  332. }
  333. }
  334. }
  335. // resolveCommand 解析启动命令:优先 node + dsh 入口 JS(无 shell 包装,避免弹窗与注入面)。
  336. func (s *Supervisor) resolveCommand() (string, []string, error) {
  337. node := s.cfg.NodePath
  338. if node == "" {
  339. node = s.nodePath()
  340. }
  341. binJS := s.cfg.DshBinJS
  342. if binJS == "" {
  343. binJS = s.dshBinJS()
  344. }
  345. port := strconv.Itoa(s.cfg.WebPort)
  346. if node != "" && binJS != "" {
  347. return node, []string{binJS, "web", "--no-open", "--host", s.cfg.WebHost, "--port", port}, nil
  348. }
  349. if p, err := exec.LookPath("dsh.cmd"); err == nil {
  350. return "cmd.exe", []string{"/c", p, "web", "--no-open", "--port", port}, nil
  351. }
  352. if p, err := exec.LookPath("dsh"); err == nil {
  353. return p, []string{"web", "--no-open", "--port", port}, nil
  354. }
  355. return "", nil, errors.New("未找到 dsh 入口:请确认已全局安装 @deepseek-ai/dsh,或在配置中指定 nodePath / dshBinJs")
  356. }