53e7d1e5f2
文档/ (2026-09-16; 老板定调: 这是 agent 底座, 所以必须扎实, 现在越扎实以后开发越简单)
00-索引 文档地图 + 30 秒概念速查 + 三条命令跑起来 + 事实源优先级
01-快速上手 体检 -> 建库 -> 起内核 -> 起驱动 -> 收工, 全带实测输出; 第一次最易踩的四个坑
02-写一个驱动 五分钟最小驱动 / 形态选择 / 能碰哪些表 / 汇报与调用两份模板 / 交付检查表
03-命令手册 两层每条命令 + 日志选项 + 退出码约定 + --json 样例 + 日常十条
04-契约与调用 一次调用的完整生命周期 / 六条仲裁 / 锁与按需拉起 / 排障表
05-日志与排障 三条道怎么读 + "症状->判据->处置"总表 + 断电收尸语义
06-架构与不变量 分层 / 14 条硬不变量 / 主流程表 / 双真相 / 为什么故意不做 / 已知薄弱点
07-模块与接口 逐模块职责与公开接口 + "想改 X -> 动哪几处"连带清单
08-数据模型 8 张表逐字段 (谁写谁读) + events.kind 字典 + 状态机 + 快照 + 排查 SQL
09-扩展指南 六个配方 (加子命令/加字段/加表/加日志来源/加自测/改判定) + 同步清单
10-验收与质量门 四道门 + 五份自测明细 + pyright 严格档 + 26 条已知坑总表 + 发布 checklist
规矩: 不重复设计文档 / 每条命令实测过再写 (含 jq 表达式) / 代码>设计>文档 的事实源优先级 /
改代码必须同步文档 (清单在 09 末尾) / 暂时没做到的事写成"已知边界"不含糊过去
修复: 同锁串行化原来是死的 (实测抓到的真缺陷)
旧行为: db.领调用 只领 pending (waiting 没人再碰) + db.同锁在跑 把 waiting 也算"占着锁"
-> 同一把锁上两条请求互相排队, 双双停在 waiting 谁也不跑 (实测 id 16/17);
而 收权超时 只收 running -> 排队连超时都没有 = 死锁
修法: ① db.领调用 的 SQL 改 state IN ('pending','waiting') -- 每轮把排队的领回来重判, 锁一空就推进
② db.同锁在跑 只认 state='running' (排队的还没拿到锁, 不挡人)
③ 内核.转发调用 waiting 分支补 deadline (排队也立期限); 内核.收权超时 遍历 running + waiting
④ 抽出 内核.期限文本() 统一算 deadline
实测: 两条同锁调用串行跑完 (19.started_at == 18.finished_at); 排队者超时被收权 (events 有记录)
回归: 自测db.py 调用组 +4 条断言 (waiting 不算占着锁 / waiting 会被重新领 / ...);
去掉一条依赖生产库全局计数的脆弱断言
其它: 内核 与 引导器 的 用法() 末尾加文档指引
验收: uvx pyright 0 errors / 0 warnings; 五份自测全过 (进程/内核 58/配置/db/日志 86);
试跑引导器.py PASS 11 / FAIL 0 / 残留无; 残留进程 0
1293 lines
58 KiB
Python
1293 lines
58 KiB
Python
#!/usr/bin/env python3
|
|
"""内核: 总调度 + 纯 CLI + PostgreSQL 当内存.
|
|
|
|
[它是什么]
|
|
四步流程 (老板 内核设计.md 那 4 行) 的落地:
|
|
① 遍历驱动目录认 配置.efi.json -> 扫描.py
|
|
② 遍历驱动配置 + 校验 -> 扫描.py
|
|
③ 生成驱动 json (注册表 + 快照) -> 扫描.py + db.py
|
|
④ 管驱动进程 (linux 命令) -> 本文件 (复用 内核/进程.py)
|
|
外加它是**总调度**: 契约匹配 / 启动排序 / 数据路由 / 级联启停 / 崩溃处理全归它.
|
|
|
|
[为什么必须有内核 (而不是让驱动互相调用)]
|
|
驱动之间零耦合: 不 import 对方, 配置里也不写对方的名字, 只声明 needs (我要什么) / provides
|
|
(我产出什么).那么"谁给谁,什么顺序,谁先谁后"就只能由唯一知道全局的一方来配 -- 内核.
|
|
驱动要别人的东西时不直连, 而是往 calls 表插一行 (want 写契约名); 内核校验后转发.这样驱动之间
|
|
连对方是谁都不知道, 自然打不起来架.
|
|
|
|
[四条定调 (老板 2026-09-15 亲定, 动手前对齐)]
|
|
1. 界面 = 纯 CLI + 日志.没有 TUI; 输出要能 `>` 重定向成文件.
|
|
2. 不设计通信协议: 驱动与内核读写同一个库, events 表就是总线.
|
|
3. 内核常驻 (甲): 运行时看着依赖链 (契约满足唤醒下游 / 上游崩级联 / 驱动请求转发仲裁).
|
|
命令走 PG commands 表 + LISTEN/NOTIFY (PG 自带的, 不是自造协议).
|
|
4. 驱动 = 一个文件夹: 里面有源码 + venv 或 exec 文件; 内核只认根目录的 配置.efi.json.
|
|
|
|
[命令一览]
|
|
(无参数) | 调度 常驻调度内核 (甲): 命令 + 调用仲裁 + 依赖巡检
|
|
列表 驱动清单 (默认动作)
|
|
扫描 只扫不启: 刷新注册表 + 快照 + 收尸
|
|
启动 <名> / 停止 <名> / 重启 <名>
|
|
状态 [名] [--json] 进程状态 (直读 /proc, 不信库里的旧 pid)
|
|
日志 <名> [-n 200] [-f] 驱动日志 (stdout/stderr 都在这一个文件里)
|
|
事件 [-n 50] 全局事件流 (events 表 = 总线)
|
|
清单 输出清单 JSON (重定向就是文件)
|
|
|
|
[判活一律回 /proc 复核]
|
|
driver_state 里的 pid 只是记账, **绝不当依据**: pid 会被系统复用, 拿旧 pid 发信号可能杀到别人
|
|
(设计 02 第 4 节三条铁律之一).所以每次动作/渲染前都调 状态.复核() -> 进程.判活().
|
|
|
|
[本版不做 (留 v0.2, 设计 02 第 9 节)]
|
|
程序/Skill 层加载,PG 角色级硬隔离 (现在是"约定 + 内核校验"),开机自启,TUI,资源限额.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import json
|
|
import os
|
|
import signal
|
|
import sys
|
|
import time
|
|
from datetime import datetime, timedelta
|
|
from pathlib import Path
|
|
from typing import Any, cast
|
|
|
|
项目根 = Path(__file__).resolve().parent.parent
|
|
# 先把 内核/ 塞进模块搜索路径, 再 import 兄弟模块 (和 UEFI.boot.py 一个套路)
|
|
sys.path.insert(0, str(项目根 / "内核"))
|
|
|
|
import db
|
|
import 扫描
|
|
import 状态
|
|
import 日志
|
|
import 进程
|
|
import 文本
|
|
|
|
版本 = "内核 v0.1"
|
|
环境文件 = 项目根 / "环境.efi.json"
|
|
内核日志路径 = 日志.内核日志路径(项目根) # 结构化日志行 (内核自己写; 引导器不再往这重定向 fd)
|
|
内核输出路径 = 日志.内核输出路径(项目根) # 进程 stdout/stderr 原始流 (命令输出 + 崩溃原文)
|
|
引导器日志路径 = 日志.引导器日志路径(项目根) # 引导器写, 内核只读 (日志 --引导器 用)
|
|
心跳秒 = 300.0 # 常驻内核多久写一条 heartbeat 事件 (断电时间线用; 太密会把总线刷满)
|
|
巡检秒 = 10.0 # 依赖链巡检间隔
|
|
调用默认时限秒 = 60.0 # 驱动调用没写 deadline 时给多久 (超时内核收权)
|
|
通知通道 = "内核" # 命令到达的唤醒通道 (LISTEN/NOTIFY 的通道名)
|
|
调度锁键 = 0x65666901 # 常驻内核的"独一份"咨询锁 (PG 会话级; 进程一死自动放, 不留死锁)
|
|
|
|
# 日志开关 (两处最容易搞错的地方之一): 门槛挡"写不写", 控制台挡"要不要同时打 stderr".
|
|
# 值从 环境.efi.json 来 (引导器写, 内核读), main() 一进来就用 用环境() 设好.
|
|
日志门槛: str = 日志.默认门槛
|
|
日志控制台: bool = True
|
|
|
|
|
|
# ─────────────────────────────── 日志与配置 ───────────────────────────────
|
|
|
|
|
|
def 说(级别: str, 消息: str) -> None:
|
|
"""写一条内核日志 (设计 02 §6): 时间 级别 [内核] 内容.
|
|
|
|
落 内核/logs/内核.log (结构化行, 只有这一种内容); 是否同时打 stderr 看 日志控制台
|
|
(守护模式关掉: 那时 stdout/stderr 都被重定向进 内核.out.log, 再从 stderr 走一遍就是副本).
|
|
写失败不拦内核干活 -- 但绝不"静默吞错": 错误内容本身照原文写出去.
|
|
"""
|
|
日志.记(内核日志路径, 级别, "内核", 消息, 门槛=日志门槛, 控制台=日志控制台)
|
|
|
|
|
|
def 用环境(环境: dict[str, Any]) -> None:
|
|
"""从环境配置装日志开关 (门槛 / 控制台).
|
|
|
|
为什么单独一个函数: 读环境() 自己也要写日志 (配置坏了就是 ERROR), 而门槛得在它之后才有--
|
|
所以 ERROR 级别的报错走默认门槛 INFO 先写出去, 装好开关再按配置过滤.
|
|
"""
|
|
global 日志门槛, 日志控制台
|
|
日志门槛 = 日志.规范化级别(环境.get("log_level") or 日志.默认门槛)
|
|
日志控制台 = 日志.控制台开()
|
|
|
|
|
|
def 读环境() -> dict[str, Any]:
|
|
"""读项目根的 环境.efi.json (内核**只读**; 写它的是引导器).
|
|
|
|
读不了 / 不是合法 JSON -> 报错退出.内核不像引导器那样能降级:
|
|
它的内存就是 PG, 而连接信息就在这份配置里, 猜不出来就别硬跑.
|
|
"""
|
|
try:
|
|
原文 = 环境文件.read_text(encoding="utf-8")
|
|
except OSError as 错:
|
|
说("ERROR", f"环境配置读不了: {错} (先跑 UEFI.boot.py --check 看环境)")
|
|
raise SystemExit(1) from 错
|
|
try:
|
|
数据: Any = json.loads(原文)
|
|
except json.JSONDecodeError as 错:
|
|
说("ERROR", f"环境配置不是合法 JSON: {错}")
|
|
raise SystemExit(1) from 错
|
|
if not isinstance(数据, dict):
|
|
说("ERROR", "环境配置顶层必须是对象 ({...})")
|
|
raise SystemExit(1)
|
|
return cast(dict[str, Any], 数据)
|
|
|
|
|
|
def 驱动根(环境: dict[str, Any]) -> Path:
|
|
"""驱动根目录 (环境.efi.json 的 driver_root; 相对路径按项目根解析)."""
|
|
值 = str(环境.get("driver_root") or "驱动")
|
|
路径 = Path(值)
|
|
return 路径 if 路径.is_absolute() else (项目根 / 路径)
|
|
|
|
|
|
def 停超时(环境: dict[str, Any]) -> float:
|
|
"""SIGTERM 之后等多久升级 SIGKILL (环境配置 stop_timeout, 默认 10s)."""
|
|
try:
|
|
return max(float(环境.get("stop_timeout") or 10), 0.5)
|
|
except (TypeError, ValueError):
|
|
return 10.0
|
|
|
|
|
|
def 日志行数(环境: dict[str, Any]) -> int:
|
|
"""日志默认显示多少行 (环境配置 log_lines, 默认 200)."""
|
|
try:
|
|
return max(int(环境.get("log_lines") or 200), 1)
|
|
except (TypeError, ValueError):
|
|
return 200
|
|
|
|
|
|
def 日志上限(环境: dict[str, Any]) -> int:
|
|
"""单份日志的轮转阈值, 单位字节 (环境配置 log_max_mb, 默认 5 MB; 0 = 不轮转)."""
|
|
原始 = 环境.get("log_max_mb")
|
|
if 原始 is None:
|
|
return 日志.默认上限字节
|
|
try:
|
|
return max(int(float(str(原始)) * 1024 * 1024), 0)
|
|
except (TypeError, ValueError):
|
|
return 日志.默认上限字节
|
|
|
|
|
|
def 日志保留(环境: dict[str, Any]) -> int:
|
|
"""轮转后留几份历史 (环境配置 log_keep, 默认 3; 0 = 不轮转)."""
|
|
try:
|
|
return max(int(环境.get("log_keep") or 日志.默认保留份数), 0)
|
|
except (TypeError, ValueError):
|
|
return 日志.默认保留份数
|
|
|
|
|
|
def 轮转日志(路径: Path, 环境: dict[str, Any], 称呼: str) -> None:
|
|
"""拉起进程之前轮转一次日志 (常驻进程的日志不能让它无限长).
|
|
|
|
为什么在启动前: 运行中的进程按 fd 追加写, 中途改名会让它继续写老 inode (等于日志丢了);
|
|
所以只在"还没起它"的时候轮转, 简单且不出错.
|
|
"""
|
|
if 日志.轮转(路径, 日志上限(环境), 日志保留(环境)):
|
|
说("WARN", f"{称呼}: 日志超过上限, 已轮转 (旧份改名 .1)")
|
|
|
|
|
|
def 连接串(环境: dict[str, Any]) -> str:
|
|
"""PG 连接串 (libpq 关键字式), 注入给驱动当环境变量 EFI_DB.
|
|
|
|
为什么走环境变量不落盘: 连接信息不落盘是老板定过的口径 (设计 01 第 3 节);
|
|
驱动拿 os.environ["EFI_DB"] 直连, 不需要问内核要地址.
|
|
"""
|
|
库 = db.从配置(环境.get("db"))
|
|
return f"host='{库.host}' port={库.port} dbname='{库.name}' user='{库.user}'"
|
|
|
|
|
|
def 字符串表(值: Any) -> list[str]:
|
|
"""把库里取出来的"字符串数组"统一成 list[str] (provides / needs / args 都是这种).
|
|
|
|
为什么要它:
|
|
`驱动.get("needs") or []` 的类型是 `Any | list[Unknown]` -- 直接迭代, 元素是 Unknown,
|
|
严格检查会满屏"类型部分未知".在这里 cast 一次, 后面到处都干净 (和引导器里 取对象/取清单
|
|
是一套思路). 顺手挡脏值: 不是数组给空表, 空串 / 非字符串项丢掉.
|
|
"""
|
|
清单 = cast(list[Any], 值) if isinstance(值, list) else []
|
|
return [文本项 for 文本项 in (str(项).strip() for 项 in 清单) if 文本项]
|
|
|
|
|
|
# ─────────────────────────────── 连库 ───────────────────────────────
|
|
|
|
|
|
def 连库(环境: dict[str, Any]) -> Any:
|
|
"""连 PG 并建表 (内核建全部 8 张, CREATE TABLE IF NOT EXISTS 幂等).
|
|
|
|
PG 连不上 -> 报错退出, **不降级到文件模式** (内存不在就没法干活; 两份真相更坏).
|
|
"""
|
|
库 = db.从配置(环境.get("db"))
|
|
try:
|
|
连接 = db.连(库)
|
|
except Exception as 错: # psycopg2 的异常类由 db.py 那层管, 这里只关心"连不上"
|
|
说("ERROR", f"连不上 PG ({库.描述()}): {错}")
|
|
raise SystemExit(1) from 错
|
|
db.建表(连接)
|
|
return 连接
|
|
|
|
|
|
def 重连(环境: dict[str, Any], 旧连接: Any) -> Any:
|
|
"""常驻内核的断线自愈: 关掉旧连接再连一次 (10 次, 每次隔 3 秒).
|
|
|
|
为什么要自愈: 常驻进程不能因为一次网络/数据库抖动就得人去敲命令重起 (跑稳了别动它的反面
|
|
是"别让它动不动就死").10 次还连不上就认输退出, 让引导器看门狗报出来.
|
|
"""
|
|
try:
|
|
旧连接.close()
|
|
except Exception:
|
|
pass
|
|
for 次 in range(1, 11):
|
|
try:
|
|
return 连库(环境)
|
|
except SystemExit:
|
|
说("WARN", f"重连 PG 失败 (第 {次}/10 次), 3 秒后再试")
|
|
time.sleep(3)
|
|
说("ERROR", "重连 PG 连续失败 10 次, 退出 (让引导器看门狗发现)")
|
|
raise SystemExit(1)
|
|
|
|
|
|
# ─────────────────────────────── 拼命令 ───────────────────────────────
|
|
|
|
|
|
def 拼命令(驱动: dict[str, Any], 环境: dict[str, Any]) -> tuple[list[str], dict[str, str], list[str]]:
|
|
"""拼出启动命令行 + 子进程环境, 返回 (argv, env, 警告表).
|
|
|
|
python 形态 (设计 01 第 2 节的解释器解析顺序):
|
|
绝对路径 -> 用它
|
|
"venv" -> <驱动根>/.venv/bin/python; 不存在 -> 回落 system 并**记警告** (不假死)
|
|
"system" -> python3
|
|
exec 形态: [入口] + args (不看后缀, 看 x 位 -- 扫描阶段已经查过 x 位)
|
|
|
|
env = 内核进程的环境 (基线) + EFI_DB (PG 连接串, 不落盘) + 驱动配置里的 env (可覆盖).
|
|
返回的警告表由调用方记进日志/事件 -- 这里不自己写库 (拼命令行这层不碰 SQL).
|
|
"""
|
|
入口 = 状态.入口路径(驱动)
|
|
args = 字符串表(驱动.get("args"))
|
|
警告表: list[str] = []
|
|
argv: list[str] = []
|
|
|
|
if str(驱动.get("runtime") or "python") == "exec":
|
|
argv = [str(入口), *args]
|
|
else:
|
|
声明 = str(驱动.get("interpreter") or "venv")
|
|
if 声明.startswith("/"):
|
|
解释器 = 声明
|
|
if not Path(解释器).exists():
|
|
警告表.append(f"配置里写的解释器不存在: {解释器} (起不起来会直接暴露)")
|
|
elif 声明 == "system":
|
|
解释器 = "python3"
|
|
else:
|
|
venv解释器 = 状态.驱动根(驱动) / ".venv" / "bin" / "python"
|
|
if venv解释器.exists():
|
|
解释器 = str(venv解释器)
|
|
else:
|
|
解释器 = "python3"
|
|
警告表.append(f"声明用 venv 但没找到 {venv解释器}, 回落 system python3")
|
|
argv = [解释器, str(入口), *args]
|
|
|
|
env = dict(os.environ)
|
|
env["EFI_DB"] = 连接串(环境)
|
|
追加 = 驱动.get("env")
|
|
if isinstance(追加, dict):
|
|
for 键, 值 in cast(dict[str, Any], 追加).items():
|
|
env[str(键)] = str(值)
|
|
return argv, env, 警告表
|
|
|
|
|
|
# ─────────────────────────────── 单驱动动作 ───────────────────────────────
|
|
|
|
|
|
def 收僵尸(pid: int) -> int | None:
|
|
"""如果这个死掉的进程是**我们自己**的孩子, wait 掉它 (不留僵尸), 顺手取出退出码.
|
|
|
|
为什么必须收:
|
|
父进程不 wait, 死掉的子进程会挂着 Z 态一直占一个 pid 名额 (停止就会永远报"还有 N 个
|
|
没收掉", 实际它早死了) -- 这个坑在引导器的进程库里踩过一次, 这里不留第二次.
|
|
只有生它的进程能 wait; 不是自己的孩子会抛 ChildProcessError, 当没这回事 (交给 init 收).
|
|
返回:
|
|
退出码 (被信号杀的是负数); 不是我们的孩子 / 还没死透 -> None.
|
|
"""
|
|
try:
|
|
结果 = os.waitpid(pid, os.WNOHANG)
|
|
except (ChildProcessError, OSError):
|
|
return None
|
|
if 结果[0] == 0:
|
|
return None
|
|
状态位 = 结果[1]
|
|
if os.WIFEXITED(状态位):
|
|
return os.WEXITSTATUS(状态位)
|
|
if os.WIFSIGNALED(状态位):
|
|
return -os.WTERMSIG(状态位)
|
|
return None
|
|
|
|
|
|
def 拉起一个(
|
|
连接: Any,
|
|
环境: dict[str, Any],
|
|
驱动: dict[str, Any],
|
|
行: dict[str, Any] | None,
|
|
理由: str = "",
|
|
) -> tuple[bool, str]:
|
|
"""把一个驱动拉起来并记账 (启动 / 重启 / 按需 / 自动重拉 全走这一份实现).
|
|
|
|
参数:
|
|
连接 / 环境: PG 与项目配置.
|
|
驱动: drivers 表一行.
|
|
行: driver_state 表一行 (用来读上一次的 restarts / pid).
|
|
理由: 非空时写进事件, 说明"为什么这次是内核主动拉的" (如 on-failure / autostart).
|
|
返回:
|
|
(起没起来, 一句话结论).
|
|
"""
|
|
名 = str(驱动.get("name"))
|
|
argv, env, 警告表 = 拼命令(驱动, 环境)
|
|
for 警告 in 警告表:
|
|
说("WARN", f"{名}: {警告}")
|
|
db.写事件(连接, "内核", "log", f"{名}: {警告}", driver=名, level="warn")
|
|
|
|
上次次数 = int(行.get("restarts") or 0) if 行 else 0
|
|
驱动日志 = 状态.日志路径(驱动)
|
|
轮转日志(驱动日志, 环境, 名)
|
|
结果 = 进程.启动(
|
|
argv,
|
|
cwd=状态.驱动根(驱动),
|
|
日志=驱动日志,
|
|
env=env,
|
|
入口=状态.入口路径(驱动),
|
|
分隔=f"启动 {名} {' '.join(argv)}",
|
|
)
|
|
|
|
if 结果.ok:
|
|
db.写状态(
|
|
连接,
|
|
名,
|
|
{
|
|
"state": 状态.运行,
|
|
"pid": 结果.pid,
|
|
"pgid": 结果.pgid,
|
|
"started_at": 状态.现在文本(),
|
|
"stopped_at": None,
|
|
"exit_code": None,
|
|
"restarts": 上次次数 + 1,
|
|
"boot_hash": str(驱动.get("config_hash") or ""),
|
|
"last_error": None,
|
|
},
|
|
)
|
|
尾巴 = f" ({理由})" if 理由 else ""
|
|
db.写事件(连接, "内核", "start", f"{名} 起来了 pid={结果.pid} pgid={结果.pgid}{尾巴}", driver=名)
|
|
return True, f"已拉起 pid={结果.pid}"
|
|
|
|
# oneshot (批处理型): 跑完就退是**正常**的 (设计 01 第 2 节), 不算失败
|
|
if str(驱动.get("mode") or "resident") == "oneshot" and 结果.exit_code == 0:
|
|
db.写状态(
|
|
连接,
|
|
名,
|
|
{
|
|
"state": 状态.已退出,
|
|
"pid": None,
|
|
"pgid": None,
|
|
"exit_code": 0,
|
|
"restarts": 上次次数 + 1,
|
|
"boot_hash": None,
|
|
"stopped_at": 状态.现在文本(),
|
|
"last_error": None,
|
|
},
|
|
)
|
|
db.写事件(连接, "内核", "exit", f"{名} oneshot 跑完就退 (退出码 0) -- 正常", driver=名)
|
|
return True, "oneshot 跑完就退 (退出码 0)"
|
|
|
|
# 其余秒退 / spawn 失败: 记 failed + 退出码 + 本次启动的日志尾巴 (不吞错)
|
|
db.写状态(
|
|
连接,
|
|
名,
|
|
{
|
|
"state": 状态.失败,
|
|
"pid": None,
|
|
"pgid": None,
|
|
"exit_code": 结果.exit_code,
|
|
"restarts": 上次次数 + 1,
|
|
"boot_hash": None,
|
|
"last_error": 结果.detail,
|
|
},
|
|
)
|
|
db.写事件(连接, "内核", "error", f"{名} 起不来: {结果.detail}", driver=名, level="error")
|
|
return False, 结果.detail
|
|
|
|
|
|
def 停一个(连接: Any, 环境: dict[str, Any], 名: str, 理由: str = "") -> tuple[bool, str]:
|
|
"""停一个驱动 (先校验 cmdline, 再 SIGTERM 进程组; 幂等).
|
|
|
|
步骤 (设计 02 第 4 节):
|
|
1. 没有 pid / /proc 里没了 -> 直接置 stopped (幂等, 不算错);
|
|
2. pid 在但 cmdline 对不上 -> 判"pid 被复用", **一个信号都不发** (宁可不杀, 不可误杀);
|
|
3. SIGTERM 给进程组 (连子树一起收), 超时升级 SIGKILL (在 进程.停止 里);
|
|
4. 停干净才置 stopped + stopped_at, 顺手收自己的僵尸孩子.
|
|
返回:
|
|
(干净没干净, 一句话结论).
|
|
"""
|
|
驱动 = db.取驱动(连接, 名)
|
|
if 驱动 is None:
|
|
return False, f"没这个驱动: {名}"
|
|
行 = db.取状态(连接, 名) or {}
|
|
pid值 = 行.get("pid")
|
|
pid = int(pid值) if isinstance(pid值, int) else None
|
|
入口 = 状态.入口路径(驱动)
|
|
活 = 进程.判活(pid, 入口)
|
|
|
|
if 活 in (进程.已停止, 进程.已崩):
|
|
db.写状态(
|
|
连接,
|
|
名,
|
|
{
|
|
"state": 状态.停止 if 活 == 进程.已停止 else 状态.崩了,
|
|
"pid": None,
|
|
"pgid": None,
|
|
"stopped_at": 状态.现在文本(),
|
|
},
|
|
)
|
|
db.写事件(连接, "内核", "stop", f"{名} 本来就不在 (幂等)", driver=名)
|
|
return True, "本来就没在跑 (幂等)"
|
|
|
|
if 活 == 进程.僵尸:
|
|
if pid is not None:
|
|
收僵尸(pid)
|
|
db.写状态(连接, 名, {"state": 状态.停止, "pid": None, "pgid": None, "stopped_at": 状态.现在文本()})
|
|
db.写事件(连接, "内核", "stop", f"{名} 已经是僵尸了 (等回收), 当它停了", driver=名)
|
|
return True, "已经是僵尸 (已死)"
|
|
|
|
if 活 == 进程.被复用 or pid is None:
|
|
db.写状态(
|
|
连接,
|
|
名,
|
|
{
|
|
"state": 状态.已退出,
|
|
"pid": None,
|
|
"pgid": None,
|
|
"stopped_at": 状态.现在文本(),
|
|
"last_error": f"pid {pid} 已被系统分给别人, 没发信号",
|
|
},
|
|
)
|
|
db.写事件(连接, "内核", "error", f"{名}: pid {pid} 被复用, 没敢杀", driver=名, level="warn")
|
|
return False, f"pid {pid} 被复用, 没敢杀 (宁可不杀, 不可误杀)"
|
|
|
|
结果 = 进程.停止(pid, 入口=入口, 超时=停超时(环境))
|
|
收僵尸(pid)
|
|
说明 = 结果.detail + (f" ({理由})" if 理由 else "")
|
|
if 结果.ok:
|
|
db.写状态(
|
|
连接,
|
|
名,
|
|
{
|
|
"state": 状态.停止,
|
|
"pid": None,
|
|
"pgid": None,
|
|
"stopped_at": 状态.现在文本(),
|
|
"last_error": None,
|
|
},
|
|
)
|
|
db.写事件(连接, "内核", "stop", f"{名} 已停: {说明}", driver=名)
|
|
return True, 说明
|
|
|
|
db.写状态(
|
|
连接,
|
|
名,
|
|
{"state": 状态.失败, "stopped_at": 状态.现在文本(), "last_error": 结果.detail},
|
|
)
|
|
db.写事件(连接, "内核", "error", f"{名} 停不干净: {结果.detail}", driver=名, level="error")
|
|
return False, 结果.detail
|
|
|
|
|
|
def 下游名(全部: dict[str, dict[str, Any]], 名: str) -> list[str]:
|
|
"""谁依赖我: needs 里出现了我 provides 的契约的驱动 (驱动之间不认识, 这层关系只有内核知道)."""
|
|
目标 = 全部.get(名)
|
|
if 目标 is None:
|
|
return []
|
|
我的契约 = set(字符串表(目标.get("provides")))
|
|
if not 我的契约:
|
|
return []
|
|
出: list[str] = []
|
|
for 别的名, 别的 in sorted(全部.items()):
|
|
if 别的名 == 名:
|
|
continue
|
|
要的 = set(字符串表(别的.get("needs")))
|
|
if 我的契约 & 要的:
|
|
出.append(别的名)
|
|
return 出
|
|
|
|
|
|
def 停序(连接: Any, 名: str, 已处理: set[str]) -> list[str]:
|
|
"""算出"停 名 的时候要按什么顺序停谁": 先递归停掉下游 (消费者), 最后停它自己.
|
|
|
|
这就是设计里的"级联生命周期": 停上游 -> 依赖它的下游一并停, 顺序是消费者先走
|
|
(别让下游在上游已经没了之后还在消费空气).
|
|
"""
|
|
if 名 in 已处理:
|
|
return []
|
|
已处理.add(名)
|
|
全部 = {str(d.get("name")): d for d in db.取全部驱动(连接)}
|
|
顺序: list[str] = []
|
|
for 下 in 下游名(全部, 名):
|
|
顺序.extend(停序(连接, 下, 已处理))
|
|
顺序.append(名)
|
|
return 顺序
|
|
|
|
|
|
# ─────────────────────────────── 驱动调用仲裁 (防打架六条) ───────────────────────────────
|
|
|
|
|
|
def 取链(行: dict[str, Any]) -> list[str]:
|
|
"""从 calls.args 里取调用链 (驱动转发时带上, 内核用它挡 A->B->A 这种环).
|
|
|
|
约定: args = {"chain": ["驱动A", "驱动B", ...]} (可省).缺省就是空链, 只靠契约成环去挡.
|
|
"""
|
|
参数 = 行.get("args")
|
|
if not isinstance(参数, dict):
|
|
return []
|
|
链 = cast(dict[str, Any], 参数).get("chain")
|
|
if not isinstance(链, list):
|
|
return []
|
|
return [str(项) for 项 in cast(list[Any], 链)]
|
|
|
|
|
|
def 校验调用(连接: Any, 环境: dict[str, Any], 契约: dict[str, str], 行: dict[str, Any]) -> tuple[str, str]:
|
|
"""给一条驱动调用请求做仲裁 (防打架六条), 返回 (判定, 说明).
|
|
|
|
判定三种:
|
|
denied 拒 (越权 / 成环 / 没人提供 / 提供方起不来)
|
|
waiting 排队 (同一把锁上已经有调用在跑)
|
|
running 放行转发
|
|
"""
|
|
调用id = int(行.get("id") or 0)
|
|
调用方 = str(行.get("caller") or "")
|
|
要的 = str(行.get("want") or "")
|
|
|
|
# 第 5 条 越权: 只能要自己 needs 里声明过的契约
|
|
调用方驱动 = db.取驱动(连接, 调用方)
|
|
if 调用方驱动 is None:
|
|
return "denied", f"调用方 {调用方} 不是注册的驱动"
|
|
声明 = set(字符串表(调用方驱动.get("needs")))
|
|
if 要的 not in 声明:
|
|
return "denied", f"越权: {调用方} 的 needs 里没有 {要的}"
|
|
|
|
# 第 2 条 成环: 契约成环在扫描期就拒载了, 这里再挡运行时链路成环.
|
|
# 链的含义是"这次请求已经走到的路" (发起方 + 它的上游), 所以只有**要的契约**或者
|
|
# **匹配出来的提供方**已经在链上才算环; 不能拿调用方自己判 -- 它一定在链尾, 那会把
|
|
# 第一次调用也拒掉 (2026-09-16 踩过: 消费器第一次请求就被判"成环").
|
|
链 = 取链(行)
|
|
提供方 = 契约.get(要的, "")
|
|
if 要的 in 链 or (提供方 != "" and 提供方 in 链):
|
|
return "denied", "调用链成环: " + " -> ".join([*链, 提供方 or 要的])
|
|
|
|
# 第 3 条 想调一个没起来的驱动: 先按需拉起, 拉不起来才拒
|
|
if not 提供方:
|
|
return "denied", f"契约没人提供: {要的}"
|
|
提供驱动 = db.取驱动(连接, 提供方)
|
|
if 提供驱动 is None:
|
|
return "denied", f"契约 {要的} 的提供方 {提供方} 不在注册表里"
|
|
提供行 = db.取状态(连接, 提供方) or {}
|
|
if 进程.判活(提供行.get("pid"), 状态.入口路径(提供驱动)) != 进程.运行中:
|
|
好, 说明 = 拉起一个(连接, 环境, 提供驱动, 提供行, f"被调用 {调用id} 按需拉起")
|
|
if not 好:
|
|
return "denied", f"提供方 {提供方} 起不来: {说明}"
|
|
|
|
# 第 1 条 同一 lock_key 串行化: 先到先执行, 后面的排队
|
|
锁 = str(行.get("lock_key") or "")
|
|
if 锁 and db.同锁在跑(连接, 锁, 调用id):
|
|
return "waiting", f"锁 {锁} 上已有调用在跑, 排队"
|
|
|
|
return "running", f"转发给 {提供方}"
|
|
|
|
|
|
def 期限文本(时限秒: float | None = None) -> str:
|
|
"""给一条调用算 deadline (ISO8601 带时区) -- 内核按它收权, 不让人无限占着.
|
|
|
|
参数:
|
|
时限秒: 多少秒之后算超时; None = 用 调用默认时限秒.
|
|
"""
|
|
秒数 = 调用默认时限秒 if 时限秒 is None else 时限秒
|
|
return (datetime.now().astimezone() + timedelta(seconds=秒数)).isoformat(timespec="seconds")
|
|
|
|
|
|
def 转发调用(连接: Any, 环境: dict[str, Any], 契约: dict[str, str], 行: dict[str, Any]) -> None:
|
|
"""仲裁 + 转发一条调用请求 (结果字段留给驱动回填: 内核只做仲裁/唤醒/收权).
|
|
|
|
转发的含义: 把 state 改成 running, 记下 provider 和 deadline, 然后 pg_notify 唤醒提供方
|
|
(常驻驱动自己 LISTEN driver_<名>).**没有任何自造协议** -- 用的都是 PG 自带的表和通知.
|
|
"""
|
|
调用id = int(行.get("id") or 0)
|
|
调用方 = str(行.get("caller") or "")
|
|
要的 = str(行.get("want") or "")
|
|
判定, 说明 = 校验调用(连接, 环境, 契约, 行)
|
|
|
|
if 判定 == "denied":
|
|
db.写调用(连接, 调用id, {"state": "denied", "error": 说明, "finished_at": 状态.现在文本()})
|
|
db.写事件(连接, "内核", "error", f"调用 {调用id} 被拒: {说明}", driver=调用方, level="warn")
|
|
return
|
|
|
|
if 判定 == "waiting":
|
|
等待改动: dict[str, Any] = {"state": "waiting", "error": 说明}
|
|
# 排队也要有期限: 第一次排队时立 deadline. 否则提供方一直占着锁, 排队者能等到天荒地老
|
|
# (2026-09-16 修: 以前 waiting 行既没 deadline 又不会被重新领 -> 排队 = 永久卡住).
|
|
if 行.get("deadline") is None:
|
|
等待改动["deadline"] = 期限文本()
|
|
db.写调用(连接, 调用id, 等待改动)
|
|
return
|
|
|
|
提供方 = 契约.get(要的, "")
|
|
db.写调用(
|
|
连接,
|
|
调用id,
|
|
{
|
|
"state": "running",
|
|
"provider": 提供方,
|
|
"started_at": 状态.现在文本(),
|
|
"deadline": 期限文本(),
|
|
"error": None,
|
|
},
|
|
)
|
|
db.写事件(
|
|
连接,
|
|
"内核",
|
|
"produce",
|
|
f"调用 {调用id}: {调用方} 要 {要的} -> 转发给 {提供方}",
|
|
driver=提供方,
|
|
)
|
|
db.通知(连接, f"driver_{提供方}", str(调用id))
|
|
|
|
|
|
def 收权超时(连接: Any) -> int:
|
|
"""给超时的调用收权 (防打架第 4 条): 标 timeout + 释放锁, 免得一个卡死的驱动拖垮全局.
|
|
|
|
两种都要收 (2026-09-16 修的):
|
|
running -- 提供方拿着锁干超时了;
|
|
waiting -- 排队的也没等到头 (提供方占着锁不松). 以前只收 running, 排队的行会永远等着.
|
|
|
|
返回:
|
|
收了几条.
|
|
"""
|
|
现在 = datetime.now().astimezone()
|
|
收 = 0
|
|
for 状态值 in ("running", "waiting"):
|
|
for 行 in db.取调用(连接, 状态值):
|
|
时限值 = 行.get("deadline")
|
|
if 时限值 is None:
|
|
continue
|
|
时刻 = 时限值 if isinstance(时限值, datetime) else None
|
|
if 时刻 is None:
|
|
continue
|
|
if 时刻.astimezone() < 现在:
|
|
db.写调用(
|
|
连接,
|
|
int(行.get("id") or 0),
|
|
{
|
|
"state": "timeout",
|
|
"error": f"超过期限 {时刻.isoformat(timespec='seconds')}, 内核收权",
|
|
"finished_at": 状态.现在文本(),
|
|
},
|
|
)
|
|
db.写事件(
|
|
连接,
|
|
"内核",
|
|
"error",
|
|
f"调用 {行.get('id')} 超时收权 (原来是 {状态值}, 提供方 {行.get('provider') or '未定'})",
|
|
driver=str(行.get("provider") or ""),
|
|
level="warn",
|
|
)
|
|
收 += 1
|
|
return 收
|
|
|
|
|
|
# ─────────────────────────────── 巡检 (常驻内核的看家活) ───────────────────────────────
|
|
|
|
|
|
def 巡检(连接: Any, 环境: dict[str, Any]) -> None:
|
|
"""周期性看一眼依赖链, 该动作的动作:
|
|
|
|
1. 收尸: 旧状态 + /proc 判现在真状态 (崩了的改库 + 记 exit 事件);
|
|
2. 按策略重拉: autostart (断电后) / restart=on-failure (自己崩的);
|
|
3. 级联: 上游不在了 -> 下游标"依赖失效" (不装作没事);
|
|
4. 调用收权: 超时的 calls 标 timeout.
|
|
"""
|
|
全部 = {str(d.get("name")): d for d in db.取全部驱动(连接)}
|
|
契约表 = 扫描.取契约(list(全部.values()))
|
|
for 名, 驱动 in 全部.items():
|
|
if not 驱动.get("valid"):
|
|
continue
|
|
行 = db.取状态(连接, 名) or {}
|
|
新状态, 原因 = 状态.复核(驱动, 行)
|
|
旧状态 = str(行.get("state") or "")
|
|
if 新状态 != 旧状态:
|
|
改动: dict[str, Any] = {"state": 新状态, "last_error": 原因 or None}
|
|
if 新状态 in (状态.已退出, 状态.崩了):
|
|
改动["pid"] = None
|
|
改动["pgid"] = None
|
|
改动["stopped_at"] = 状态.现在文本()
|
|
死pid = 行.get("pid")
|
|
if isinstance(死pid, int):
|
|
退出码 = 收僵尸(死pid)
|
|
if 退出码 is not None:
|
|
改动["exit_code"] = 退出码
|
|
db.写状态(连接, 名, 改动)
|
|
db.写事件(
|
|
连接,
|
|
"内核",
|
|
"exit",
|
|
f"{名}: {旧状态} -> {新状态} ({原因})" if 原因 else f"{名}: {旧状态} -> {新状态}",
|
|
driver=名,
|
|
level="warn",
|
|
)
|
|
行 = db.取状态(连接, 名) or 行
|
|
理由 = 状态.该拉起(驱动, 行, 新状态)
|
|
if 理由:
|
|
拉起一个(连接, 环境, 驱动, 行, 理由)
|
|
# 级联: 上游没在跑 -> 下游标"依赖失效" (状态还是 running, 但把话说清楚)
|
|
if str(行.get("state") or "") == 状态.运行:
|
|
缺: list[str] = []
|
|
for 契约名 in 字符串表(驱动.get("needs")):
|
|
上 = 契约表.get(契约名, "")
|
|
if not 上 or 上 == 名:
|
|
continue
|
|
上行 = db.取状态(连接, 上) or {}
|
|
if 进程.判活(上行.get("pid"), 状态.入口路径(全部.get(上, {}))) == 进程.运行中:
|
|
continue
|
|
# 只有"上游崩了 / 起不来"才算下游依赖失效; 上游只是还没启动 (等着按需拉起)
|
|
# 属于正常状态, 不能天天报假警 (2026-09-16 踩过: 没启动也标"依赖失效").
|
|
上状态 = str(上行.get("state") or 状态.停止)
|
|
if 上状态 in (状态.崩了, 状态.失败):
|
|
缺.append(f"{契约名}(提供方 {上}, 状态 {状态.显示状态(上状态)})")
|
|
if 缺:
|
|
说明 = "依赖失效: 上游不在跑 " + ", ".join(缺)
|
|
if str(行.get("last_error") or "") != 说明:
|
|
db.写状态(连接, 名, {"last_error": 说明})
|
|
db.写事件(连接, "内核", "error", f"{名}: {说明}", driver=名, level="warn")
|
|
收权超时(连接)
|
|
|
|
|
|
# ─────────────────────────────── CLI 动作 ───────────────────────────────
|
|
|
|
|
|
def 打印列表(结果: 扫描.扫描结果, 环境: dict[str, Any]) -> None:
|
|
"""打印驱动清单表 + 页脚 (设计 02 第 6 节的输出样例)."""
|
|
行表: list[list[str]] = []
|
|
状态表 = {str(行.get("name")): 行 for 行 in 结果.状态表}
|
|
驱动表 = {str(d.get("name")): d for d in 结果.驱动}
|
|
for 名 in sorted(驱动表):
|
|
驱动 = 驱动表[名]
|
|
行 = 状态表.get(名) or {}
|
|
if not 驱动.get("valid"):
|
|
行表.append([状态.显示状态(状态.无效), 名, str(驱动.get("runtime") or ""), "—", "—", "—",
|
|
str(驱动.get("error") or "")])
|
|
continue
|
|
活, 原因 = 状态.复核(驱动, 行)
|
|
pid值 = 行.get("pid")
|
|
信息 = 进程.读进程信息(int(pid值)) if isinstance(pid值, int) else None
|
|
时长 = 信息.存活文本 if (活 == 状态.运行 and 信息 is not None) else "—"
|
|
配置列 = "待重启" if 状态.待重启(驱动, 行) else "一致"
|
|
说明 = 原因 or (str(驱动.get("note") or "") or "—")
|
|
if 活 == 状态.运行 and 状态.待重启(驱动, 行):
|
|
说明 = "配置已改, 需重启生效"
|
|
行表.append([
|
|
状态.显示状态(活),
|
|
名,
|
|
str(驱动.get("runtime") or ""),
|
|
str(pid值) if isinstance(pid值, int) and 活 == 状态.运行 else "—",
|
|
时长,
|
|
配置列,
|
|
说明,
|
|
])
|
|
if not 行表:
|
|
行表.append(["—", "(一个驱动都没有)", "—", "—", "—", "—", "驱动根: " + str(驱动根(环境))])
|
|
文本.打印(文本.表格(["状态", "驱动名", "形态", "PID", "运行时长", "配置", "说明"], 行表))
|
|
print()
|
|
print(f"驱动 {结果.总数} / 有效 {结果.有效} / 运行 {结果.在跑} / 无效 {结果.无效}"
|
|
f" 扫描 #{结果.清单版本} {datetime.now().astimezone().strftime('%H:%M:%S')}")
|
|
for 警告 in 结果.问题:
|
|
print(f" [WARN] {警告}")
|
|
|
|
|
|
def 命令列表(连接: Any, 环境: dict[str, Any]) -> int:
|
|
"""列表: 扫一遍再渲染 (列表要反映磁盘现状, 不能拿旧注册表糊弄)."""
|
|
结果 = 扫描.扫描(连接, 驱动根(环境), 版本)
|
|
打印列表(结果, 环境)
|
|
return 0
|
|
|
|
|
|
def 命令扫描(连接: Any, 环境: dict[str, Any]) -> int:
|
|
"""扫描: 只扫不启 (刷新注册表 + 快照 + 收尸), 打印一行汇总.
|
|
|
|
"只扫不启"是有意的: 想启动谁就显式敲 `启动 <名>`; autostart 的拉起发生在常驻内核启动时.
|
|
"""
|
|
结果 = 扫描.扫描(连接, 驱动根(环境), 版本)
|
|
print(f"扫描完成: 驱动 {结果.总数} / 有效 {结果.有效} / 运行 {结果.在跑} / 无效 {结果.无效}"
|
|
f" 扫描 #{结果.清单版本}")
|
|
for 名, 理由 in 结果.该拉起:
|
|
print(f" [建议] {名}: {理由} (敲 启动 {名} 或让常驻内核去拉)")
|
|
for 警告 in 结果.问题:
|
|
print(f" [WARN] {警告}")
|
|
return 0
|
|
|
|
|
|
def 命令启动(连接: Any, 环境: dict[str, Any], 名: str) -> int:
|
|
"""启动一个驱动 (幂等: 已经在跑就直说, 不重复起)."""
|
|
驱动 = db.取驱动(连接, 名)
|
|
if 驱动 is None:
|
|
print(f"没这个驱动: {名} (先 扫描 看看驱动根里有什么)")
|
|
return 1
|
|
if not 驱动.get("valid"):
|
|
print(f"驱动 {名} 配置不合法, 不能启动: {驱动.get('error')}")
|
|
return 1
|
|
行 = db.取状态(连接, 名) or {}
|
|
入口 = 状态.入口路径(驱动)
|
|
if 进程.判活(行.get("pid"), 入口) == 进程.运行中:
|
|
print(f"{名} 已经在跑 (pid {行.get('pid')}) -- 幂等, 不重复起")
|
|
return 0
|
|
好, 说明 = 拉起一个(连接, 环境, 驱动, 行)
|
|
print(f"{名}: {说明}")
|
|
return 0 if 好 else 1
|
|
|
|
|
|
def 命令停止(连接: Any, 环境: dict[str, Any], 名: str) -> int:
|
|
"""停一个驱动 (级联: 依赖它的下游先停)."""
|
|
驱动 = db.取驱动(连接, 名)
|
|
if 驱动 is None:
|
|
print(f"没这个驱动: {名}")
|
|
return 1
|
|
顺序 = 停序(连接, 名, set())
|
|
码 = 0
|
|
for 目标 in 顺序:
|
|
理由 = "" if 目标 == 名 else f"上游 {名} 要停"
|
|
好, 说明 = 停一个(连接, 环境, 目标, 理由)
|
|
print(f"{目标}: {说明}")
|
|
if not 好:
|
|
码 = 1
|
|
return 码
|
|
|
|
|
|
def 命令重启(连接: Any, 环境: dict[str, Any], 名: str) -> int:
|
|
"""重启一个驱动: 先停干净 (含级联), 再起 (一次只动一件事)."""
|
|
停码 = 命令停止(连接, 环境, 名)
|
|
if 停码 != 0:
|
|
print("停得不干净, 先不起了 (别在没停干净的进程上再叠一个)")
|
|
return 停码
|
|
行 = db.取状态(连接, 名) or {}
|
|
if 进程.判活(行.get("pid"), 状态.入口路径(db.取驱动(连接, 名) or {})) == 进程.运行中:
|
|
print("还有进程活着, 先不起了")
|
|
return 1
|
|
return 命令启动(连接, 环境, 名)
|
|
|
|
|
|
def 命令状态(连接: Any, 环境: dict[str, Any], 参数: list[str]) -> int:
|
|
"""状态 [名] [--json]: 不给名字就出整张清单 (跟 列表 一样), 给了就出单个详情.
|
|
|
|
判活走 状态.复核() -> 进程.判活() (直读 /proc/<pid>/cmdline, 不信库里的旧 pid);
|
|
跟库里的状态不一致时**如实标出来**, 但不改库 (这是只读命令, 收尸是 扫描/巡检 的活).
|
|
"""
|
|
要json = "--json" in 参数
|
|
参数 = [项 for 项 in 参数 if not 项.startswith("--")]
|
|
if not 参数:
|
|
if 要json:
|
|
return 命令清单(连接, 环境)
|
|
结果 = 扫描.扫描(连接, 驱动根(环境), 版本)
|
|
打印列表(结果, 环境)
|
|
return 0
|
|
|
|
名 = 参数[0]
|
|
驱动 = db.取驱动(连接, 名)
|
|
if 驱动 is None:
|
|
print(f"没这个驱动: {名}")
|
|
return 1
|
|
行 = db.取状态(连接, 名) or {}
|
|
活, 原因 = 状态.复核(驱动, 行)
|
|
pid值 = 行.get("pid")
|
|
信息 = 进程.读进程信息(int(pid值)) if isinstance(pid值, int) else None
|
|
if 要json:
|
|
出 = {
|
|
"efi": 1,
|
|
"kernel": 版本,
|
|
"ts": 状态.现在文本(),
|
|
"注册表": 驱动,
|
|
"状态": 状态.组装快照(驱动, 行, 版本),
|
|
"现场": {
|
|
"判定": 活,
|
|
"原因": 原因,
|
|
"存活秒": 信息.存活秒 if 信息 is not None else None,
|
|
"内存KB": 信息.rss_kb if 信息 is not None else None,
|
|
"cmdline": 信息.cmdline if 信息 is not None else [],
|
|
},
|
|
}
|
|
print(json.dumps(出, ensure_ascii=False, indent=2, default=str))
|
|
return 0
|
|
|
|
行表 = [
|
|
["项", "值"],
|
|
["状态", 状态.显示状态(活) + (f" ({原因})" if 原因 else "")],
|
|
["形态", f"{驱动.get('runtime')} mode={驱动.get('mode')} restart={驱动.get('restart')}"],
|
|
["入口", str(状态.入口路径(驱动))],
|
|
["PID / 进程组", f"{pid值} / {行.get('pgid')}" if pid值 else "—"],
|
|
["启动时间", str(行.get("started_at") or "—")],
|
|
["运行时长", 信息.存活文本 if 信息 is not None else "—"],
|
|
["内存 / CPU", f"{信息.内存文本} / {信息.cpu秒:.1f}s" if 信息 is not None else "—"],
|
|
["命令行", 信息.cmdline文本 if 信息 is not None else "—"],
|
|
["配置", ("待重启 (配置改过, 重启才生效)" if 状态.待重启(驱动, 行) else "一致")],
|
|
["重启次数", str(行.get("restarts") or 0)],
|
|
["契约", "提供: " + (", ".join(字符串表(驱动.get("provides"))) or "—")
|
|
+ " 需要: " + (", ".join(字符串表(驱动.get("needs"))) or "—")],
|
|
["上次错误", str(行.get("last_error") or "—")],
|
|
]
|
|
文本.打印(文本.表格(行表[0], 行表[1:], [文本.左, 文本.左]))
|
|
if 原因:
|
|
print(f"\n 注意: 现判 {状态.显示状态(活)} 而库里记的是 {状态.显示状态(行.get('state'))} -- "
|
|
f"库里的状态是上次记账, 以现场判定为准")
|
|
return 0
|
|
|
|
|
|
def 命令日志(连接: Any, 环境: dict[str, Any], 参数: list[str]) -> int:
|
|
"""日志 [驱动名] [-n 200] [-f] [--级别 X] [-g 关键词] [--json] [--全部] [--内核] [--引导器]
|
|
|
|
四种形态 (同一个日志子系统, 区别只是"看谁的"):
|
|
日志 驱动日志台账 (谁有日志 / 多大 / 几行 / 最后改动)
|
|
日志 <驱动名> ... 某个驱动的日志 (默认尾部 200 行)
|
|
日志 --全部 ... 所有驱动的日志尾部 (每条前一行 == 驱动名 ==)
|
|
日志 --内核 / --引导器 内核自己的结构化日志 / 引导器的动作日志
|
|
|
|
驱动日志是**别人的原始 stdout** (内核只重定向不解析): 所以级别只能靠字样猜 (猜级别),
|
|
内核/引导器日志是本库写的固定格式, 级别是真字段 (--json 出去也是拆好的字段).
|
|
"""
|
|
选, 问题 = 日志.解析选项(参数, 日志行数(环境))
|
|
if 问题:
|
|
print(f"日志参数有问题: {问题}")
|
|
print(用法日志())
|
|
return 1
|
|
if 选.跟随 and (选.全部 or 选.内核 or 选.引导器 or 选.输出):
|
|
print("-f 一次只能跟一个来源 (驱动名 / --内核 / --引导器 / --输出)")
|
|
return 1
|
|
if 选.输出:
|
|
return 打印某个日志(内核输出路径, 选, "内核输出")
|
|
if 选.内核:
|
|
return 打印某个日志(内核日志路径, 选, "内核")
|
|
if 选.引导器:
|
|
return 打印某个日志(引导器日志路径, 选, "引导器")
|
|
if 选.全部:
|
|
return 打印全部驱动日志(连接, 选)
|
|
if not 选.名:
|
|
return 打印日志台账(连接)
|
|
驱动 = db.取驱动(连接, 选.名)
|
|
if 驱动 is None:
|
|
print(f"没这个驱动: {选.名}")
|
|
return 1
|
|
return 打印某个日志(状态.日志路径(驱动), 选, 选.名)
|
|
|
|
|
|
def 打印某个日志(路径: Path, 选: 日志.选项, 称呼: str) -> int:
|
|
"""打一个日志文件 (带过滤选项); 内核/引导器/驱动共用这一份实现."""
|
|
if 选.跟随:
|
|
print(f"[跟 {称呼}] {路径} (Ctrl+C 停)", flush=True)
|
|
return 日志.跟(路径, 选.级别, 选.关键词)
|
|
行表 = 日志.尾(路径, 选.行数, 选.级别, 选.关键词)
|
|
if not 行表:
|
|
if not 路径.exists():
|
|
print(f"还没有 {称呼} 的日志文件: {路径}")
|
|
else:
|
|
print(f"{路径} 里没有符合条件的行 (级别 {选.级别 or '全部'} / 关键词 {选.关键词 or '无'})")
|
|
return 0
|
|
日志.打印(行表, 选.json输出)
|
|
return 0
|
|
|
|
|
|
def 打印全部驱动日志(连接: Any, 选: 日志.选项) -> int:
|
|
"""`日志 --全部`: 每个驱动一段, 带 == 名字 == 分隔 (一眼看全, 不用一个个敲)."""
|
|
驱动表 = db.取全部驱动(连接)
|
|
if not 驱动表:
|
|
print("还没有驱动")
|
|
return 0
|
|
没日志 = 0
|
|
for 驱动 in 驱动表:
|
|
名 = str(驱动.get("name"))
|
|
路径 = 状态.日志路径(驱动)
|
|
print(f"\n== {名} ({路径}) ==")
|
|
行表 = 日志.尾(路径, 选.行数, 选.级别, 选.关键词)
|
|
if not 行表:
|
|
print(" (没有日志 / 没有符合条件的行)")
|
|
没日志 += 1
|
|
continue
|
|
日志.打印(行表, 选.json输出)
|
|
print(f"\n共 {len(驱动表)} 个驱动, {没日志} 个没有日志")
|
|
return 0
|
|
|
|
|
|
def 打印日志台账(连接: Any) -> int:
|
|
"""日志台账: 先看这张表, 再决定跟谁的日志 (以前只能靠猜哪个驱动有日志)."""
|
|
驱动表 = db.取全部驱动(连接)
|
|
if not 驱动表:
|
|
print("还没有驱动")
|
|
return 0
|
|
行表: list[list[str]] = []
|
|
for 驱动 in 驱动表:
|
|
路径 = 状态.日志路径(驱动)
|
|
文件表 = 日志.轮转名单(路径)
|
|
总量 = 0
|
|
for 文件 in 文件表:
|
|
try:
|
|
总量 += 文件.stat().st_size
|
|
except OSError:
|
|
continue
|
|
行数 = len(日志.读全部(路径))
|
|
最后 = "—"
|
|
try:
|
|
最后 = datetime.fromtimestamp(路径.stat().st_mtime).astimezone().isoformat(timespec="seconds")
|
|
except OSError:
|
|
最后 = "—"
|
|
历史 = f" (+{len(文件表) - 1} 份历史)" if len(文件表) > 1 else ""
|
|
行表.append([
|
|
str(驱动.get("name")),
|
|
"有" if 文件表 else "无",
|
|
日志.大小文本(总量),
|
|
str(行数),
|
|
f"{最后}{历史}",
|
|
])
|
|
文本.打印(文本.表格(["驱动", "日志", "占用", "行数", "最后改动"], 行表))
|
|
print("\n提示: 日志 <驱动名> [-n 200] [-f] [--级别 WARN] [-g 关键词] [--json]")
|
|
return 0
|
|
|
|
|
|
def 用法日志() -> str:
|
|
"""日志命令的用法一行 (参数写错时打)."""
|
|
return ("用法: 日志 [驱动名] [-n 200] [-f] [--级别 DEBUG|INFO|WARN|ERROR] "
|
|
"[-g 关键词] [--json] [--全部] [--内核] [--引导器] [--输出]")
|
|
|
|
|
|
def 命令事件(连接: Any, 环境: dict[str, Any], 参数: list[str]) -> int:
|
|
"""事件 [-n 50]: 全局事件流 (events 表 = 总线; 内核写, 驱动也写)."""
|
|
条数 = 50
|
|
if "-n" in 参数:
|
|
位置 = 参数.index("-n")
|
|
if 位置 + 1 < len(参数):
|
|
try:
|
|
条数 = max(int(参数[位置 + 1]), 1)
|
|
except ValueError:
|
|
print(f"-n 后面要跟一个数字, 收到 {参数[位置 + 1]!r}")
|
|
return 1
|
|
_ = 环境
|
|
事件表 = db.读事件(连接, 条数)
|
|
行表: list[list[str]] = []
|
|
for 行 in 事件表:
|
|
时刻 = 行.get("ts")
|
|
时刻文本 = 时刻.strftime("%m-%d %H:%M:%S") if isinstance(时刻, datetime) else str(时刻 or "")
|
|
行表.append([
|
|
时刻文本,
|
|
str(行.get("level") or "info"),
|
|
str(行.get("kind") or ""),
|
|
str(行.get("source") or ""),
|
|
str(行.get("driver") or "—"),
|
|
str(行.get("message") or ""),
|
|
])
|
|
if not 行表:
|
|
print("还没有事件")
|
|
return 0
|
|
文本.打印(文本.表格(["时间", "级别", "类型", "来源", "驱动", "内容"], 行表))
|
|
print(f"\n共 {len(行表)} 条 (最近的在上面)")
|
|
return 0
|
|
|
|
|
|
def 命令清单(连接: Any, 环境: dict[str, Any]) -> int:
|
|
"""清单: 输出 JSON (重定向就是文件; 机器读的, 不混人读的表格)."""
|
|
状态表 = {str(行.get("name")): 行 for 行 in db.取全部状态(连接)}
|
|
驱动表 = db.取全部驱动(连接)
|
|
出 = {
|
|
"efi": 1,
|
|
"kernel": 版本,
|
|
"ts": 状态.现在文本(),
|
|
"driver_root": str(驱动根(环境)),
|
|
"驱动": [
|
|
{
|
|
"注册表": 驱动,
|
|
"状态": 状态.组装快照(驱动, 状态表.get(str(驱动.get("name"))), 版本),
|
|
}
|
|
for 驱动 in 驱动表
|
|
],
|
|
}
|
|
print(json.dumps(出, ensure_ascii=False, indent=2, default=str))
|
|
return 0
|
|
|
|
|
|
# ─────────────────────────────── 常驻调度 (甲) ───────────────────────────────
|
|
|
|
|
|
def 执行命令(连接: Any, 环境: dict[str, Any], 命令: str, 参数: list[str]) -> tuple[int, str]:
|
|
"""按命令名执行一条命令, 返回 (退出码, 一句话结果).
|
|
|
|
CLI 和常驻内核**共用这一份实现** (只此一份, 不复制): CLI 是短命的执行者,
|
|
常驻内核是排队命令的执行者, 干的是同一件事.
|
|
"""
|
|
if 命令 in ("列表", "") and not 参数:
|
|
return 命令列表(连接, 环境), "已输出驱动清单"
|
|
if 命令 == "扫描":
|
|
return 命令扫描(连接, 环境), "已扫描"
|
|
if 命令 == "启动" and 参数:
|
|
return 命令启动(连接, 环境, 参数[0]), f"启动 {参数[0]}"
|
|
if 命令 == "停止" and 参数:
|
|
return 命令停止(连接, 环境, 参数[0]), f"停止 {参数[0]}"
|
|
if 命令 == "重启" and 参数:
|
|
return 命令重启(连接, 环境, 参数[0]), f"重启 {参数[0]}"
|
|
if 命令 == "状态":
|
|
return 命令状态(连接, 环境, 参数), "已输出状态"
|
|
if 命令 == "日志":
|
|
return 命令日志(连接, 环境, 参数), f"日志 {' '.join(参数)}".strip()
|
|
if 命令 == "事件":
|
|
return 命令事件(连接, 环境, 参数), "已输出事件"
|
|
if 命令 == "清单":
|
|
return 命令清单(连接, 环境), "已输出清单"
|
|
return 2, f"不认识这个命令: {命令} {参数}"
|
|
|
|
|
|
def 命令调度(连接: Any, 环境: dict[str, Any]) -> int:
|
|
"""常驻调度内核 (甲): 命令消费 + 驱动调用仲裁 + 依赖巡检 + 心跳.
|
|
|
|
为什么必须常驻 (设计 02 第 4 节): 总调度要在**运行时**看着依赖链 -- 契约满足要唤醒下游,
|
|
上游崩要级联,驱动请求要转发仲裁, 只在启动那一下算一遍不够.
|
|
|
|
退出时**不顺手停驱动**: 驱动进程是独立会话, 内核死了它们照跑; 下次启动时收尸逻辑会把它们
|
|
认领回来 (设计 03 说的"内核可以随时死"就是这个意思).要停驱动请显式敲 停止.
|
|
"""
|
|
停止标志 = {"停": False}
|
|
|
|
def 收工(号: int, 帧: object) -> None:
|
|
"""收到 TERM/INT: 只立个旗子, 让主循环自己收尾 (信号处理函数里别干重活)."""
|
|
停止标志["停"] = True
|
|
说("INFO", f"收到信号 {号}, 准备退出 (驱动不动)")
|
|
|
|
# 独一份: 抢 PG 会话级咨询锁, 抢不到说明已经有一个内核在跑 (两个调度器=打架)
|
|
if not db.试锁(连接, 调度锁键):
|
|
说("ERROR", "已经有一个内核在跑了 (抢不到调度锁) -- 同一份 PG 内存只能有一个调度器, 别起第二个")
|
|
return 1
|
|
signal.signal(signal.SIGTERM, 收工)
|
|
signal.signal(signal.SIGINT, 收工)
|
|
|
|
轮转日志(内核日志路径, 环境, "内核") # 常驻内核前先轮转, 别让它无限长 (改 fd 之前做)
|
|
说("INFO", f"{版本} 常驻调度启动 pid={os.getpid()}")
|
|
收了 = db.收尸命令(连接, "内核上次没跑完 (断电 / 被杀)")
|
|
if 收了:
|
|
说("WARN", f"收尸: 补了 {收了} 条没写完的命令")
|
|
结果 = 扫描.扫描(连接, 驱动根(环境), 版本)
|
|
打印列表(结果, 环境)
|
|
db.写事件(连接, "内核", "start", f"{版本} 常驻调度启动 pid={os.getpid()}")
|
|
for 名, 理由 in 结果.该拉起:
|
|
说("INFO", f"autostart: {名} -- {理由}")
|
|
拉起一个(连接, 环境, db.取驱动(连接, 名) or {}, db.取状态(连接, 名), 理由)
|
|
db.监听(连接, 通知通道)
|
|
|
|
上次心跳 = 0.0
|
|
上次巡检 = 0.0
|
|
上次扫描 = time.monotonic()
|
|
while not 停止标志["停"]:
|
|
try:
|
|
连接.poll() # psycopg2 自带: 把网上的 NOTIFY 收进内存 (不阻塞)
|
|
通知表 = db.收通知(连接)
|
|
for 通道, 载荷 in 通知表:
|
|
if 通道 == 通知通道:
|
|
说("INFO", f"收到命令唤醒 (id={载荷})")
|
|
|
|
# 1) 领命令执行 (CLI 直接执行过的不在这里: 它们的 state 不是 pending)
|
|
命令行 = db.领命令(连接)
|
|
while 命令行 is not None:
|
|
名 = str(命令行.get("cmd") or "")
|
|
参数原始 = 命令行.get("args")
|
|
参数: list[str] = []
|
|
if isinstance(参数原始, dict):
|
|
值 = cast(dict[str, Any], 参数原始).get("参数")
|
|
if isinstance(值, list):
|
|
参数 = [str(项) for 项 in cast(list[Any], 值)]
|
|
说("INFO", f"执行排队命令 #{命令行.get('id')}: {名} {参数}")
|
|
码, 说明 = 执行命令(连接, 环境, 名, 参数)
|
|
db.记命令结果(连接, int(命令行.get("id") or 0), "done" if 码 == 0 else "failed",
|
|
{"码": 码, "说明": 说明})
|
|
命令行 = db.领命令(连接)
|
|
|
|
# 2) 领调用请求: 仲裁 + 转发 (驱动只写表, 不跟内核握手)
|
|
契约表 = 扫描.取契约(db.取全部驱动(连接))
|
|
for 调用 in db.领调用(连接):
|
|
转发调用(连接, 环境, 契约表, 调用)
|
|
|
|
# 3) 巡检 (每 巡检秒): 收尸 / 重拉 / 级联 / 超时收权
|
|
现在 = time.monotonic()
|
|
if 现在 - 上次巡检 >= 巡检秒:
|
|
巡检(连接, 环境)
|
|
上次巡检 = 现在
|
|
|
|
# 4) 心跳 (每 心跳秒): 断电时间线用 (库里最后一条 heartbeat 就是死的那一刻)
|
|
if 现在 - 上次心跳 >= 心跳秒:
|
|
db.写事件(连接, "内核", "heartbeat", f"{版本} 常驻中 (pid {os.getpid()})",
|
|
data={"驱动数": len(db.取全部驱动(连接))})
|
|
上次心跳 = 现在
|
|
|
|
# 5) 定期重扫 (每 10 分钟): 磁盘上新加了驱动/改了配置, 常驻内核也得看见
|
|
if 现在 - 上次扫描 >= 600:
|
|
扫描.扫描(连接, 驱动根(环境), 版本)
|
|
上次扫描 = 现在
|
|
|
|
time.sleep(0.5)
|
|
except SystemExit:
|
|
raise
|
|
except Exception as 错: # 单轮出错不能让常驻内核死: 记下来, 重连, 继续
|
|
说("WARN", f"这一轮出错 (会继续跑): {错}")
|
|
time.sleep(2)
|
|
连接 = 重连(环境, 连接)
|
|
|
|
说("INFO", "常驻调度退出 (驱动没动, 下次启动会认领它们)")
|
|
try:
|
|
db.写事件(连接, "内核", "exit", f"{版本} 常驻调度退出 (收到停止信号)")
|
|
except Exception as 错:
|
|
说("WARN", f"退出事件没写进去 (PG 可能已经断了): {错}")
|
|
return 0
|
|
|
|
|
|
# ─────────────────────────────── 入口 ───────────────────────────────
|
|
|
|
|
|
def 用法() -> str:
|
|
"""命令一览 (没参数/敲错的时候打给老板看)."""
|
|
return "\n".join([
|
|
"用法: python3 内核/内核.py <命令>",
|
|
"",
|
|
" (无参数) | 调度 常驻调度内核 (甲): 命令 + 调用仲裁 + 依赖巡检 (引导器用它守护)",
|
|
" 列表 驱动清单 (默认动作)",
|
|
" 扫描 只扫不启: 刷新注册表 + 快照 + 收尸",
|
|
" 启动 <驱动名> 拉起一个驱动 (幂等: 在跑就直说)",
|
|
" 停止 <驱动名> 停一个驱动 (先校验 cmdline 再 SIGTERM 进程组; 级联停下游)",
|
|
" 重启 <驱动名> 停干净了才启 (一次只动一件事)",
|
|
" 状态 [驱动名] [--json] 进程状态 (直读 /proc, 不信库里的旧 pid)",
|
|
" 日志 [驱动名] [选项] 日志 (不写名字 = 台账; --全部 / --内核 / --引导器)",
|
|
" -n 200 尾部行数",
|
|
" -f 实时跟 (先吐 10 行, Ctrl+C 停)",
|
|
" --级别 WARN 只要这个级别及以上 (驱动日志按字样猜级别)",
|
|
" -g 关键词 只要含它的行",
|
|
" --json 输出 JSON Lines (机器读; 结构化日志拆成字段)",
|
|
" 事件 [-n 50] 全局事件流 (events 表 = 总线)",
|
|
" 清单 输出清单 JSON (重定向就是文件)",
|
|
"",
|
|
" 判活一律回 /proc 复核; 停之前先校验 cmdline (pid 会被系统复用, 宁可不杀不可误杀)",
|
|
" 日志分三条道: 内核.log(结构化) / 内核.out.log(命令输出+崩溃原文) / 驱动/<名>/logs(驱动原始输出)",
|
|
" 文档: 文档/00-索引.md (用起来 01-05 / 改底座 06-10; 设计档案在 设计/)",
|
|
])
|
|
|
|
|
|
def main(argv: list[str]) -> int:
|
|
"""CLI 入口: 短命进程执行一条命令; 无参数 = 常驻调度.
|
|
|
|
每条命令都往 commands 表落一行 (审计): CLI 自己执行时用 state='running' 落行,
|
|
这样常驻内核**不会**把 CLI 已经干过的活再领一遍 (它只领 pending).
|
|
"""
|
|
环境 = 读环境()
|
|
用环境(环境) # 装日志开关 (门槛/控制台); 读环境自己报错时用的是默认门槛 INFO
|
|
连接 = 连库(环境)
|
|
try:
|
|
if not argv or argv[0] in ("调度", "常驻"):
|
|
return 命令调度(连接, 环境)
|
|
|
|
命令, *参数 = argv
|
|
# 每条命令都落一行审计 (state=running: 我自己在执行, 别让常驻内核再领一遍)
|
|
命令id = db.记命令(连接, "cli", 命令, {"参数": 参数}, state="running")
|
|
码, 说明 = 执行命令(连接, 环境, 命令, 参数)
|
|
if 命令id:
|
|
db.记命令结果(连接, 命令id, "done" if 码 == 0 else "failed", {"码": 码, "说明": 说明})
|
|
if 码 == 2:
|
|
print(用法())
|
|
print(f"\n[错误] {说明}")
|
|
return 码
|
|
finally:
|
|
连接.close()
|
|
|
|
|
|
if __name__ == "__main__":
|
|
sys.exit(main(sys.argv[1:]))
|