Files
efi-kernel/内核/db.py
T
lou 53e7d1e5f2 底座文档 11 份 + 修掉"同锁调用互相排队、双双卡死"
文档/ (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
2026-09-16 21:13:31 +08:00

906 lines
37 KiB
Python

"""唯一碰 SQL 的文件: 连接 / 建表 / 事件 / 台账.
[为什么只有这一个文件写 SQL (设计 02 §3 定死)]
SQL 散在全项目 = 表结构一改就得满地找.集中在这里, 谁想读写库就调这里的函数.
[谁用哪部分]
* 内核 -> drivers / driver_state / events / scans / commands / calls (驱动全生命周期)
* 引导器 -> kernel_env (体检台账) / kernel_runs (内核运行台账), 它只建自己这两张,
不碰驱动程序的地盘 (分层: 引导器管内核, 内核管驱动)
* 驱动 -> 只写 events / calls (约定 + 内核校验; v0.2 再上 PG 角色做硬隔离)
[连接从哪来]
项目根的 环境.efi.json 的 db 段: {name, host, port, user}.
host 填的是 **unix socket 目录** (如 /home/lou/pgdata/socket), 不走 TCP -- 本机库没必要绕网络栈.
驱动侧的连接串由内核经环境变量 EFI_DB 注入, 不落盘 (密钥不落盘口径).
[降级约定 (很重要)]
psycopg2 是**可选**依赖:
* 引导器 缺它 -> 只 WARN, 降级成"只写 json 快照", 不算失败 (引导器必须永远能跑起来报错)
* 内核 缺它 -> 直接报错退出 (它的内存就是 PG, 没有内存没法干活, 不许静默降级)
`有psycopg2` 这个开关就是给调用方判这个的.
"""
from __future__ import annotations
import json
import re
from dataclasses import dataclass
from typing import Any
try:
# 这两行带 pyright: ignore[reportMissingModuleSource]:
# psycopg2 是编译版轮子 (只有 .so 没有 .py 源码), pyright 因此报"无法从源码解析导入" --
# 对 C 扩展模块这是误报. 类型来自 .venv 里装的 psycopg2-stubs (PEP 561, 见 环境.efi.json).
import psycopg2 as _pg # pyright: ignore[reportMissingModuleSource]
import psycopg2.extras as _pgx # pyright: ignore[reportMissingModuleSource]
except ImportError: # 引导器零第三方依赖, 缺了就降级
_pg = None
_pgx = None
# 给调用方判"要不要走降级分支"
有psycopg2: bool = _pg is not None
# ─────────────────────────────── 建表语句 ───────────────────────────────
# 全部 CREATE TABLE IF NOT EXISTS: 幂等, 跑多少次都一样, 所以内核每次 boot 都可以直接调 建表().
# 每句都带 IF NOT EXISTS, _建() 靠它回一个"建/查了哪张表"的清单给人看.
# 内核的表 (8 张: 司机 + 状态 + 总线 + 扫描批次 + 命令 + 调用请求 + 引导器两张台账)
内核表: list[str] = [
# ① 驱动注册表: 内核扫描后 upsert.
# config_hash 是 配置.efi.json 的 sha256 -- 跟 driver_state.boot_hash 一比就知道
# "配置改过但还没重启" (不用 diff 内容).valid/error 是配置校验的结果, 隔离失败用.
"""CREATE TABLE IF NOT EXISTS drivers (
name text PRIMARY KEY,
dir text NOT NULL,
runtime text NOT NULL,
entry text NOT NULL,
interpreter text,
args jsonb DEFAULT '[]',
env jsonb DEFAULT '{}',
deps text[] DEFAULT '{}',
provides text[] DEFAULT '{}',
needs text[] DEFAULT '{}',
autostart boolean DEFAULT false,
restart text DEFAULT 'no',
mode text DEFAULT 'resident',
config_hash text,
entry_hash text,
valid boolean DEFAULT true,
error text,
note text,
scanned_at timestamptz DEFAULT now()
)""",
# ② 运行时状态: 一行一驱动, 内核每次动作刷新.pid 只是记录, 判活一律回 /proc 复核
# (pid 会被系统复用).boot_hash = 起进程那一刻的配置指纹.
"""CREATE TABLE IF NOT EXISTS driver_state (
name text PRIMARY KEY REFERENCES drivers(name) ON DELETE CASCADE,
state text NOT NULL DEFAULT 'stopped',
pid integer,
pgid integer,
started_at timestamptz,
stopped_at timestamptz,
exit_code integer,
restarts integer DEFAULT 0,
boot_hash text,
list_version bigint,
last_error text,
updated_at timestamptz DEFAULT now()
)""",
# ③ 事件流 = 总线: 内核写, 驱动也写.kind 约定 start|stop|exit|log|produce|heartbeat|error.
# level 约定 info|warn|error.data 是 jsonb, 装不确定的载荷 (表结构就是消息格式, 不另设协议).
"""CREATE TABLE IF NOT EXISTS events (
id bigserial PRIMARY KEY,
ts timestamptz DEFAULT now(),
source text NOT NULL,
driver text,
level text DEFAULT 'info',
kind text,
message text,
data jsonb
)""",
"CREATE INDEX IF NOT EXISTS events_ts_idx ON events (ts DESC)",
"CREATE INDEX IF NOT EXISTS events_driver_idx ON events (driver, ts DESC)",
# ④ 扫描批次: list_version 自增, 用来判"这份状态是不是本次扫描的".
"""CREATE TABLE IF NOT EXISTS scans (
list_version bigserial PRIMARY KEY,
started_at timestamptz DEFAULT now(),
kernel text,
total int,
valid int,
invalid int,
running int
)""",
# ⑤ 命令表: 内核常驻(甲)才需要 -- CLI 客户端把命令写进来, 常驻内核 LISTEN/NOTIFY 消费.
"""CREATE TABLE IF NOT EXISTS commands (
id bigserial PRIMARY KEY,
ts timestamptz DEFAULT now(),
source text,
cmd text NOT NULL,
args jsonb DEFAULT '{}',
state text DEFAULT 'pending',
result jsonb,
started_at timestamptz,
finished_at timestamptz
)""",
# ⑥ 调用请求: 驱动要别人的产出时写这里 (want 填**契约名**, 不填驱动名).
# provider 是内核自己记的账, 请求方看不到 -- 所以驱动之间永远不认识彼此.
# lock_key 防同时写同一份数据; deadline 是超时线 (内核收权用); 调用链成环会被 denied.
"""CREATE TABLE IF NOT EXISTS calls (
id bigserial PRIMARY KEY,
ts timestamptz DEFAULT now(),
caller text NOT NULL,
want text NOT NULL,
args jsonb DEFAULT '{}',
state text DEFAULT 'pending',
provider text,
lock_key text,
result jsonb,
error text,
deadline timestamptz,
started_at timestamptz,
finished_at timestamptz
)""",
]
# 引导器的两张台账 (设计 03 §7).引导器只建这两张, 内核建全部.
引导器表: list[str] = [
# 环境/包体检台账: 每次 自检/环境/透传 之前落一行, 留痕给"环境什么时候变坏的"用.
# packages 存 jsonb 数组: [{name, want, got, ok, required}]
"""CREATE TABLE IF NOT EXISTS kernel_env (
id bigserial PRIMARY KEY,
ts timestamptz DEFAULT now(),
python_version text,
venv_path text,
venv_healthy boolean,
packages jsonb,
pg_ok boolean,
driver_root_ok boolean,
ok boolean,
detail text
)""",
# 内核运行台账: 每次运行一行.finished_at IS NULL = 还没收尾 (进程在跑, 或者断电/被杀留下的).
# mode: oneshot (前台跑一次) | daemon (后台常驻).detail 存失败原因 + 日志尾巴.
"""CREATE TABLE IF NOT EXISTS kernel_runs (
id bigserial PRIMARY KEY,
started_at timestamptz DEFAULT now(),
finished_at timestamptz,
argv text,
mode text,
pid integer,
exit_code integer,
seconds real,
ok boolean,
detail text
)""",
]
# ─────────────────────────────── 连接 ───────────────────────────────
@dataclass
class 数据库:
"""一份 PG 连接信息 (来自 环境.efi.json 的 db 段).
字段:
name: 库名 (本项目 = efi_kernel).
host: **unix socket 目录** (不是主机名/IP), 如 /home/lou/pgdata/socket.
port: 端口 (5432); socket 文件名就是 .s.PGSQL.<port>.
user: 系统用户名 (PG 用 peer 认证, 所以写本机用户名).
"""
name: str
host: str
port: int
user: str
def 连接参数(self, 库名: str = "") -> dict[str, Any]:
"""psycopg2.connect(**参数) 用的字典.库名传空 = 连 self.name."""
return {
"dbname": 库名 or self.name,
"user": self.user,
"host": self.host,
"port": self.port,
}
def 描述(self) -> str:
"""人读的一行连接描述 (报错,日志里贴这个, 别只写"连不上")."""
return f"{self.name} @ {self.host}:{self.port} (user={self.user})"
def 从配置(: dict[str, Any] | None) -> 数据库:
"""把 环境.efi.json 的 db 段转成 数据库.
参数:
段: 配置里的 db 字典; None 或字段缺失都给保守默认 (不让缺配置直接崩).
返回:
数据库 实例.
默认值的来路:
优先吃环境变量 PGHOST / PGPORT / PGUSER, 再兜底到本机常见位置 (~/pgdata/socket, 5432, 当前用户).
这样同一个项目挪到别的机器上, 改环境变量或改配置都行, 不用动代码.
"""
: dict[str, Any] = if is not None else {}
端口原始 = .get("port", 5432)
try:
端口 = int(端口原始)
except (TypeError, ValueError):
端口 = 5432
return 数据库(
name=str(.get("name", "efi_kernel")),
host=str(.get("host", "/home/lou/pgdata/socket")),
port=端口,
user=str(.get("user", "lou")),
)
def _游标工厂() -> Any:
"""字典游标工厂 (RealDictCursor).
为什么统一用它:
默认游标 fetch 出来是 tuple, 取列得按下标 (行[0]) -- 加一列就全乱.
字典游标能按列名取 (行["id"]), 表结构改了也不容易错 (踩过: RETURNING id 用默认游标报
"tuple indices must be integers or slices, not str").
"""
if _pgx is None:
raise RuntimeError("没装 psycopg2 -- 这一层不可用 (引导器请走降级分支)")
return _pgx.RealDictCursor
def (: 数据库, 超时: float = 5.0, 库名: str = "") -> Any:
"""开一条连接.连不上就抛异常, 由调用方决定是降级还是退出 (本层不做决定).
参数:
库: 连接信息.
超时: 连接超时秒 (下限 1s).
库名: 传了就改连这个库 (如查库存在时先连 postgres 管理库).
返回:
psycopg2 连接 (autocommit = True).
为什么 autocommit:
台账是一行一落的小写操作, 不该被调用方的事务卡住 (体检落库失败也不该回滚掉别的).
真需要事务的复杂流程 (内核的调度) 自己显式开事务.
"""
if _pg is None:
raise RuntimeError("没装 psycopg2")
参数 = dict(.连接参数(库名))
参数["connect_timeout"] = max(int(超时), 1)
连接 = _pg.connect(**参数)
连接.autocommit = True
return 连接
# ─────────────────────────────── 建表 ───────────────────────────────
def 建表(连接: Any) -> list[str]:
"""内核用: 八张表全建 (幂等).返回"建/查了哪些对象"的清单."""
return _建(连接, 内核表 + 引导器表)
def 建引导器表(连接: Any) -> list[str]:
"""引导器用: 只动它自己那两张 (不碰驱动程序的地盘, 分层不越界)."""
return _建(连接, 引导器表)
def _建(连接: Any, 语句表: list[str]) -> list[str]:
"""逐句执行 DDL; 从句子里的 IF NOT EXISTS 后面抠出对象名回给调用方.
返回:
对象名列表 (表名 / 索引名), 顺序跟语句一致.执行出错直接抛 (建表失败不能装看不见).
"""
: list[str] = []
with 连接.cursor() as 游标:
for 语句 in 语句表:
对象 = re.search(r"IF NOT EXISTS\s+(\S+)", 语句)
游标.execute(语句)
.append(对象.group(1) if 对象 is not None else 语句[:40])
return
def 库存在(连接: Any, 库名: str) -> bool:
"""查这个库在不在 (引导器体检第 6 项用: 实例通了但库没建是最常见的半通状态).
参数:
连接: 建议连到 postgres 管理库再查 (目标库不存在时连不上去).
库名: 要查的库名.
返回:
存在 True / 不存在 False.
"""
with 连接.cursor() as 游标:
游标.execute("SELECT 1 FROM pg_database WHERE datname = %s", (库名,))
return 游标.fetchone() is not None
# ─────────────────────────────── 事件 (总线) ───────────────────────────────
def 写事件(
连接: Any,
source: str,
kind: str,
message: str,
driver: str | None = None,
level: str = "info",
data: dict[str, Any] | None = None,
) -> None:
"""往 events 表插一行 -- 这是"没有通信协议"的落地方式: 谁想汇报就插一行, 内核不用解析 stdout.
参数:
连接: PG 连接.
source: 谁写的 (kernel / 引导器 / <驱动名>).
kind: 事件类型约定 start|stop|exit|log|produce|heartbeat|error.
message: 人读的一句话.
driver: 关联的驱动名 (可空).
level: info|warn|error.
data: 结构化载荷 (jsonb), 不确定的内容放这里.
返回:
None.写失败会抛异常 (事件丢了不该静默).
"""
with 连接.cursor() as 游标:
游标.execute(
"INSERT INTO events (source, driver, level, kind, message, data)"
" VALUES (%s, %s, %s, %s, %s, %s::jsonb)",
(source, driver, level, kind, message, json.dumps(data or {}, ensure_ascii=False)),
)
# ─────────────── 内核: 注册表 / 运行时状态 / 扫描批次 / 命令 / 调用 ───────────────
# 内核三块地盘的读写全在这儿 (唯一碰 SQL 的文件, 别处一律不写 SQL).
# 约定: text[] 列传 list, jsonb 列传 json.dumps 出来的串 + ::jsonb 强转.
def 记驱动(连接: Any, 驱动: dict[str, Any]) -> None:
"""把一个驱动的扫描结果 upsert 进注册表 (drivers 表).
参数:
连接: PG 连接.
驱动: 扫描出来的记录, 键名跟 drivers 表的列名一致.
必需: name / dir / runtime / entry;
可选: interpreter / args / env / provides / needs / autostart / restart / mode /
config_hash / entry_hash / valid / error / note.
返回:
无.
说明:
驱动名是主键, 所以 ON CONFLICT(name) 覆盖 -- 每次扫描都以磁盘为准刷一遍.
valid=false 的记录**照样入库** (带 error 原因): 列表要能看见"这个驱动坏了",
而不是假装它不存在 (假成功比报错更坏).
"""
with 连接.cursor() as 游标:
游标.execute(
"INSERT INTO drivers ("
" name, dir, runtime, entry, interpreter, args, env, provides, needs,"
" autostart, restart, mode, config_hash, entry_hash, valid, error, note, scanned_at)"
" VALUES (%s, %s, %s, %s, %s, %s::jsonb, %s::jsonb, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, now())"
" ON CONFLICT (name) DO UPDATE SET"
" dir = EXCLUDED.dir, runtime = EXCLUDED.runtime, entry = EXCLUDED.entry,"
" interpreter = EXCLUDED.interpreter, args = EXCLUDED.args, env = EXCLUDED.env,"
" provides = EXCLUDED.provides, needs = EXCLUDED.needs, autostart = EXCLUDED.autostart,"
" restart = EXCLUDED.restart, mode = EXCLUDED.mode, config_hash = EXCLUDED.config_hash,"
" entry_hash = EXCLUDED.entry_hash, valid = EXCLUDED.valid, error = EXCLUDED.error,"
" note = EXCLUDED.note, scanned_at = now()",
(
str(驱动.get("name", "")),
str(驱动.get("dir", "")),
str(驱动.get("runtime", "python")),
str(驱动.get("entry", "")),
驱动.get("interpreter"),
json.dumps(驱动.get("args") or [], ensure_ascii=False),
json.dumps(驱动.get("env") or {}, ensure_ascii=False),
list(驱动.get("provides") or []),
list(驱动.get("needs") or []),
bool(驱动.get("autostart", False)),
str(驱动.get("restart", "no")),
str(驱动.get("mode", "resident")),
驱动.get("config_hash"),
驱动.get("entry_hash"),
bool(驱动.get("valid", True)),
驱动.get("error"),
驱动.get("note"),
),
)
def 清不在(连接: Any, 保留: list[str]) -> list[str]:
"""把"磁盘上已经没有"的驱动从注册表删掉, 返回被删的名字表.
为什么要删:
文件夹没了但注册表留着 -> 列表里冒出幽灵驱动, 启动它必然失败.
driver_state 是 ON DELETE CASCADE, 状态行跟着走, 不留孤儿状态.
注意:
只删 drivers, 不碰 events -- 事件是历史, 历史不删.
"""
with 连接.cursor(cursor_factory=_游标工厂()) as 游标:
游标.execute("DELETE FROM drivers WHERE NOT (name = ANY(%s)) RETURNING name", (保留,))
return [str(["name"]) for in 游标.fetchall()]
def 取驱动(连接: Any, name: str) -> dict[str, Any] | None:
"""按名字取一个驱动 (注册表); 取不到返回 None."""
= _查(连接, "SELECT * FROM drivers WHERE name = %s", (name,))
return [0] if else None
def 取全部驱动(连接: Any) -> list[dict[str, Any]]:
"""取全部驱动, 按名字排序 (列表 / 扫描 / 契约匹配都要用)."""
return _查(连接, "SELECT * FROM drivers ORDER BY name", ())
# driver_state 允许被 写状态() 改动的列 (动态 UPDATE 的白名单, 挡住拼错的列名)
状态列: set[str] = {
"state",
"pid",
"pgid",
"started_at",
"stopped_at",
"exit_code",
"restarts",
"boot_hash",
"list_version",
"last_error",
}
def 确保状态行(连接: Any, name: str, 清单版本: int | None = None) -> None:
"""保证 driver_state 里有这一行 (没有就插一行 stopped).
扫描时对每个驱动调一次; 之后内核的所有改动都走 UPDATE (不用每次都想"要不要 INSERT").
"""
with 连接.cursor() as 游标:
游标.execute(
"INSERT INTO driver_state (name, state, list_version, updated_at)"
" VALUES (%s, 'stopped', %s, now()) ON CONFLICT (name) DO NOTHING",
(name, 清单版本),
)
def 写状态(连接: Any, name: str, 改动: dict[str, Any]) -> None:
"""改 driver_state 的若干列 (只改传进来的列, 其它列一个都不动).
为什么动态拼 SQL:
内核一次动作只改一部分 (启动改 pid/pgid/started_at; 停止改 state/exit_code/stopped_at).
写死列会逼调用方把整行传一遍, 传漏就把别人的字段抹成 NULL -- 状态就是这么丢的.
护栏:
列名必须在 状态列 白名单里, 不在就抛 ValueError; 值一律走 %s 参数, 不拼进 SQL.
"""
= sorted(k for k in 改动 if k not in 状态列)
if :
raise ValueError(f"driver_state 没有这些列: {}")
= list(改动.keys())
if not :
return
片段 = ", ".join(f"{k} = %s" for k in )
= tuple(改动[k] for k in )
with 连接.cursor() as 游标:
游标.execute(f"UPDATE driver_state SET {片段}, updated_at = now() WHERE name = %s", (*, name))
def 取状态(连接: Any, name: str) -> dict[str, Any] | None:
"""按名字取一个驱动的运行时状态; 取不到返回 None."""
= _查(连接, "SELECT * FROM driver_state WHERE name = %s", (name,))
return [0] if else None
def 取全部状态(连接: Any) -> list[dict[str, Any]]:
"""取全部运行时状态, 按名字排序 (列表 / 收尸巡检用)."""
return _查(连接, "SELECT * FROM driver_state ORDER BY name", ())
def 记扫描批次(连接: Any, 内核版本: str, 总数: int, 有效: int, 无效: int, 在跑: int) -> int:
"""记一次扫描, 返回本批次的 list_version (自增主键).
list_version 的用处: 判断"手里这份状态是不是本次扫描的" -- 旧版本的状态行一眼能认出来.
"""
with 连接.cursor(cursor_factory=_游标工厂()) as 游标:
游标.execute(
"INSERT INTO scans (kernel, total, valid, invalid, running)"
" VALUES (%s, %s, %s, %s, %s) RETURNING list_version",
(内核版本, 总数, 有效, 无效, 在跑),
)
= 游标.fetchone()
return int(["list_version"]) if is not None else 0
def 取最近扫描(连接: Any) -> dict[str, Any] | None:
"""取最近一次扫描批次 (列表页脚要显示 扫描 #N 时间); 一次都没扫过返回 None."""
= _查(
连接,
"SELECT list_version, started_at, kernel, total, valid, invalid, running"
" FROM scans ORDER BY list_version DESC LIMIT 1",
(),
)
return [0] if else None
def 读事件(连接: Any, 条数: int = 50, driver: str | None = None) -> list[dict[str, Any]]:
"""读事件流 (新的在前). driver 给了就只看这个驱动的."""
if driver:
return _查(
连接,
"SELECT * FROM events WHERE driver = %s ORDER BY ts DESC, id DESC LIMIT %s",
(driver, 条数),
)
return _查(连接, "SELECT * FROM events ORDER BY ts DESC, id DESC LIMIT %s", (条数,))
def 记命令(
连接: Any,
source: str,
cmd: str,
args: dict[str, Any] | None = None,
state: str = "pending",
) -> int:
"""往 commands 表写一条命令 (审计 + 给常驻内核排队), 返回命令 id.
state='pending' = 交给常驻内核去执行 (脚本/自动化的用法);
CLI 自己执行的命令用 state='running' 落行, 免得被常驻内核重复领走.
"""
with 连接.cursor(cursor_factory=_游标工厂()) as 游标:
游标.execute(
"INSERT INTO commands (source, cmd, args, state, started_at)"
" VALUES (%s, %s, %s::jsonb, %s, CASE WHEN %s = 'pending' THEN NULL ELSE now() END)"
" RETURNING id",
(source, cmd, json.dumps(args or {}, ensure_ascii=False), state, state),
)
= 游标.fetchone()
return int(["id"]) if is not None else 0
def 领命令(连接: Any) -> dict[str, Any] | None:
"""常驻内核领一条待执行的命令 (pending -> running), 没有就返回 None.
为什么用 FOR UPDATE SKIP LOCKED:
万一有两个内核进程 (手滑起了两份), 它们不会领到同一条命令 -- 一条命令只被执行一次.
"""
with 连接.cursor(cursor_factory=_游标工厂()) as 游标:
游标.execute(
"UPDATE commands SET state = 'running', started_at = now()"
" WHERE id = (SELECT id FROM commands WHERE state = 'pending'"
" ORDER BY id LIMIT 1 FOR UPDATE SKIP LOCKED)"
" RETURNING *"
)
= 游标.fetchone()
return dict() if is not None else None
def 记命令结果(连接: Any, 命令id: int, state: str, result: dict[str, Any] | None = None) -> None:
"""回填命令结果 (state: done | failed), CLI 轮询这个字段取结果."""
with 连接.cursor() as 游标:
游标.execute(
"UPDATE commands SET state = %s, result = %s::jsonb, finished_at = now() WHERE id = %s",
(state, json.dumps(result or {}, ensure_ascii=False), 命令id),
)
def 读命令(连接: Any, 命令id: int) -> dict[str, Any] | None:
"""读一条命令的当前状态 (CLI 轮询等结果用)."""
= _查(连接, "SELECT * FROM commands WHERE id = %s", (命令id,))
return [0] if else None
# calls 允许被 写调用() 改动的列 (同 状态列 的护栏思路)
调用列: set[str] = {
"state",
"provider",
"lock_key",
"result",
"error",
"deadline",
"started_at",
"finished_at",
}
def 领调用(连接: Any) -> list[dict[str, Any]]:
"""常驻内核领走所有待处理的驱动调用请求 (pending / waiting -> waiting), 返回这些行.
驱动只写 calls 表, 不跟内核握手 -- 内核轮询这一张表就等于收请求 (没有协议).
为什么 pending 和 waiting 都要领 (2026-09-16 实测踩到):
只领 pending 时, "因为锁被别人占着"而排队的行会**永远卡在 waiting** -- 没人再碰它,
而 收权超时 又只收 running, 于是排队 = 死。更糟的是两条同锁请求同批进来会**互相排队**:
甲看乙在 waiting、乙看甲在 waiting, 双双不动 (实测: 两条都停在 waiting, 谁也不跑)。
所以每轮把 waiting 一起领回来重判: 锁一空就自动推进成 running。
判重的开销可以忽略 (表小, 一轮 0.5s), 换来的是"排队真能排到头".
注意 SQL 写法: UPDATE 不能直接带 ORDER BY (PG 会报 syntax error at or near "ORDER");
要排序 + 跳锁就得用子查询 (跟 领命令 同一个套路, 这里踩过一次).
"""
with 连接.cursor(cursor_factory=_游标工厂()) as 游标:
游标.execute(
"UPDATE calls SET state = 'waiting'"
" WHERE id IN (SELECT id FROM calls WHERE state IN ('pending', 'waiting')"
" ORDER BY id FOR UPDATE SKIP LOCKED)"
" RETURNING *"
)
return [dict() for in 游标.fetchall()]
def 写调用(连接: Any, 调用id: int, 改动: dict[str, Any]) -> None:
"""改 calls 的若干列 (state/provider/lock_key/result/error/deadline/...); 列名走白名单.
result 是 jsonb 列: 给 dict / list 自动 dumps + ::jsonb 强转 (跟 记命令结果 同一个口径),
给字符串就按原样进. 之前不转 -- 直接塞 dict 会被 psycopg2 顶回来
("can't adapt type 'dict'"), 2026-09-16 由 内核/自测db.py 抓出来 (当时内核只写
state/error/deadline, 没踩到, 但调用方迟早会拿它写 result).
"""
= sorted(k for k in 改动 if k not in 调用列)
if :
raise ValueError(f"calls 没有这些列: {}")
= list(改动.keys())
if not :
return
片段 = ", ".join(f"{k} = %s" + ("::jsonb" if k == "result" else "") for k in )
: list[Any] = []
for k in :
v = 改动[k]
.append(json.dumps(v, ensure_ascii=False) if k == "result" and isinstance(v, (dict, list)) else v)
with 连接.cursor() as 游标:
游标.execute(f"UPDATE calls SET {片段} WHERE id = %s", (*, 调用id))
def 读调用(连接: Any, 调用id: int) -> dict[str, Any] | None:
"""读一条调用请求 (转发后查它的状态)."""
= _查(连接, "SELECT * FROM calls WHERE id = %s", (调用id,))
return [0] if else None
def 取调用(连接: Any, state: str, 条数: int = 200) -> list[dict[str, Any]]:
"""按状态取一批调用请求 (内核巡检用: 找超时的 running, 找排队的 waiting)."""
return _查(
连接,
"SELECT * FROM calls WHERE state = %s ORDER BY id LIMIT %s",
(state, 条数),
)
def 同锁在跑(连接: Any, lock_key: str, 排除id: int) -> bool:
"""同一个 lock_key 上是不是**已经有调用真正在跑** (防打架第 1 条: 同一份数据串行化).
只认 state='running' (2026-09-16 修的): 把 'waiting' 也算进"占着锁"会让两条同锁请求
**互相排队** -- 甲看乙 waiting 就排队、乙看甲 waiting 也排队, 两条一起卡死 (实测过).
排队的行应该看成"还没拿到锁", 所以它不挡别人; 真正挡人的是拿到锁在干活那条.
"""
= _查(
连接,
"SELECT id FROM calls WHERE lock_key = %s AND state = 'running' AND id <> %s LIMIT 1",
(lock_key, 排除id),
)
return bool()
def 监听(连接: Any, 通道: str) -> None:
"""LISTEN 一个通道 (内核用它等"有命令/有调用"的唤醒信号).
通道名在 SQL 里是标识符, 不能参数化 -> 先按白名单校验 (只允许字母/数字/下划线/中文) 再拼.
"""
if not 通道 or not all(字符.isalnum() or 字符 == "_" for 字符 in 通道):
raise ValueError(f"通道名不合法: {通道!r}")
with 连接.cursor() as 游标:
游标.execute(f"LISTEN {通道}")
def 收通知(连接: Any) -> list[tuple[str, str]]:
"""把 PG 攒着的通知全取出来, 返回 (通道, 载荷) 列表.
用法: 每轮先 连接.poll() (psycopg2 自带的, 把网络上的通知收进内存), 再调这个取走.
"""
: list[tuple[str, str]] = []
while 连接.notifies:
通知 = 连接.notifies.pop(0)
.append((str(通知.channel), str(通知.payload or "")))
return
def 通知(连接: Any, 通道: str, 载荷: str) -> None:
"""pg_notify 发一条通知: 只是把睡着的常驻进程叫醒, 真实数据在表里 (通知不落盘).
这是 PG 自带的机制, 不是自造协议 (老板定调第 2 条: 不设计通信协议).
约定: 驱动侧的通道名是 driver_<驱动名>, 所以驱动名别用连字符/空格 (那会让 LISTEN 的
标识符写法失效); 一律用字母 / 数字 / 中文 / 下划线.
"""
with 连接.cursor() as 游标:
游标.execute("SELECT pg_notify(%s, %s)", (通道, 载荷))
def 试锁(连接: Any, : int) -> bool:
"""抢一把 PG **会话级**咨询锁: 拿到 True / 别人拿着 False.
用途: 常驻调度内核必须是**独一份** -- 两个调度器同时巡检 / 同时拉驱动就是打架,
而且它们改的是同一份 PG 内存 (设计 02 第 3 节: 内存只有一份, 谁都不能自作主张).
会话级锁的妙处: 进程一死连接就断, 锁自动放掉 -- 断电 / 被杀都不会留下死锁,
下一次启动照样能起来 (不需要手工清锁).
"""
= _查(连接, "SELECT pg_try_advisory_lock(%s) AS 拿到", (,))
return bool([0]["拿到"]) if else False
def 收尸命令(连接: Any, 详情: str) -> int:
"""把"上一次内核被断电/杀掉时没写完"的命令收尾 (state running -> failed), 返回条数.
跟引导器收 kernel_runs 是一个道理: 断电留下的 running 行会永远挂着, 没人认领.
"""
with 连接.cursor() as 游标:
游标.execute(
"UPDATE commands SET state = 'failed', result = %s::jsonb, finished_at = now()"
" WHERE state = 'running' AND finished_at IS NULL",
(json.dumps({"detail": 详情}, ensure_ascii=False),),
)
return int(游标.rowcount)
# ─────────────────────────────── 引导器两张台账 ───────────────────────────────
def 记体检(连接: Any, 快照: dict[str, Any]) -> None:
"""把一份体检快照写进 kernel_env (引导器 自检/环境/透传 时调).
参数:
连接: PG 连接.
快照: 就是写进 环境状态.efi.json 的那份 dict (同一份数据两处落盘: 文件 + 库).
返回:
None.取值走 取bool() 容错 (快照是 json, 值可能是 None 或字符串).
"""
= 快照.get("packages", [])
with 连接.cursor() as 游标:
游标.execute(
"INSERT INTO kernel_env"
" (python_version, venv_path, venv_healthy, packages, pg_ok, driver_root_ok, ok, detail)"
" VALUES (%s, %s, %s, %s::jsonb, %s, %s, %s, %s)",
(
str(快照.get("python", {}).get("version", "")),
str(快照.get("venv", {}).get("path", "")),
取bool(快照.get("venv", {}).get("healthy")),
json.dumps(, ensure_ascii=False),
取bool(快照.get("pg", {}).get("ok")),
取bool(快照.get("driver_root", {}).get("ok")),
取bool(快照.get("blocking_ok")),
str(快照.get("detail", "")),
),
)
def 取bool(: Any) -> bool:
"""把 json 里的值统一成 bool.
为什么需要:
快照是 json 序列化过的, 字段可能是 None / 字符串 "true" / 真 bool.
PG 的 boolean 列不认字符串, 先在这边归好再插.
"""
if is None:
return False
if is True or is False:
return
return str().strip().lower() in ("1", "true", "yes", "ok")
def 记运行开始(连接: Any, argv: str, mode: str, pid: int | None) -> int:
"""内核要起进程了: 先插一行 (finished_at 留空 = 还在跑), 返回行 id 给后面收尾用.
参数:
连接: PG 连接.
argv: 完整命令行原文 (出事时能看出到底跑的什么).
mode: oneshot (前台跑一次) | daemon (后台常驻).
pid: 常驻进程的 pid; 前台透传还没起时给 None.
返回:
kernel_runs.id; 拿不到返回 0 (调用方按 0 当"没记成").
"""
工厂 = _游标工厂()
with 连接.cursor(cursor_factory=工厂) as 游标:
游标.execute(
"INSERT INTO kernel_runs (argv, mode, pid) VALUES (%s, %s, %s) RETURNING id",
(argv, mode, pid),
)
= 游标.fetchone()
if is None:
return 0
return int(["id"])
def 记运行结束(
连接: Any,
运行id: int,
exit_code: int | None,
seconds: float,
ok: bool,
detail: str = "",
) -> None:
"""给某次运行收尾: 补 finished_at / 退出码 / 耗时 / ok / detail.
参数:
运行id: 记运行开始() 返回的 id.
exit_code: 退出码 (None = 没拿到).
seconds: 耗时秒.
ok: 成功与否.
detail: 一句话结论, 失败时贴日志尾巴 (不吞错).
"""
with 连接.cursor() as 游标:
游标.execute(
"UPDATE kernel_runs SET finished_at = now(), exit_code = %s, seconds = %s, ok = %s, detail = %s"
" WHERE id = %s",
(exit_code, seconds, ok, detail, 运行id),
)
def 记内核pid(连接: Any, 运行id: int, pid: int | None) -> None:
"""补上常驻内核的 pid (启动成功后才知道, 所以单独一次 UPDATE)."""
with 连接.cursor() as 游标:
游标.execute("UPDATE kernel_runs SET pid = %s WHERE id = %s", (pid, 运行id))
def 最近运行(连接: Any) -> dict[str, Any] | None:
"""最后一次运行记录 (内核 状态 用: 上次什么时候跑的,多久,退出码,ok 还是没标)."""
行表 = _查(连接, "SELECT * FROM kernel_runs ORDER BY id DESC LIMIT 1", ())
if not 行表:
return None
return 行表[0]
def 未结束运行(连接: Any) -> list[dict[str, Any]]:
"""所有 finished_at 为空的记录 (断电/被杀留下的"假运行中").
用途: 内核 状态 提示"PG 里有 N 条没收尾", 内核 停止 时统一收尾.
"""
return _查(连接, "SELECT * FROM kernel_runs WHERE finished_at IS NULL ORDER BY id", ())
def 收尾未结束(连接: Any, detail: str) -> int:
"""把所有没写 finished_at 的行补上 (标 ok=false + 写明是谁收的尾).
参数:
连接: PG 连接.
detail: 收尾原因 (如 "引导器停止内核时收尾").
返回:
补了几行 (调用方打印出来, 让老板知道有没有断电残留).
"""
with 连接.cursor() as 游标:
游标.execute(
"UPDATE kernel_runs SET finished_at = now(), ok = false, detail = %s"
" WHERE finished_at IS NULL",
(detail,),
)
return int(游标.rowcount)
def 今日运行次数(连接: Any) -> int:
"""今天 (PG 所在时区的当天 0 点起) 跑了几次内核."""
行表 = _查(
连接,
"SELECT count(*) AS 次数 FROM kernel_runs WHERE started_at >= date_trunc('day', now())",
(),
)
if not 行表:
return 0
return int(行表[0]["次数"])
def (连接: Any, sql: str, 参数: tuple[Any, ...] = ()) -> list[dict[str, Any]]:
"""公开的只读查询 (自测 / 诊断用; 生产流程请走上面那些有名字的函数).
存在的唯一理由: 自测要按下标核对**真落库**的值 (行数 / 某一列), 而 _查 是私有的 --
跨模块碰私有符号会被类型检查判 reportPrivateUsage. 这一步纯转手, 没有任何额外逻辑.
"""
return _查(连接, sql, 参数)
def _查(连接: Any, sql: str, 参数: tuple[Any, ...]) -> list[dict[str, Any]]:
"""内部: 跑一条 SELECT 并把结果转成 dict 列表 (字典游标 + dict() 一次到位).
参数:
连接: PG 连接.
sql: 带 %s 占位符的语句 (**永不用字符串拼参数**, 走驱动层的参数绑定).
参数: 占位符的值.
返回:
[{列名: 值, ...}, ...]; 没结果返回空表.
"""
工厂 = _游标工厂()
with 连接.cursor(cursor_factory=工厂) as 游标:
游标.execute(sql, 参数)
行表 = 游标.fetchall()
: list[dict[str, Any]] = []
for in 行表:
.append(dict())
return