插件与自注册后端设计(v0.1 草案)¶
目标:用户不改 taskmq 源码就能接自己的后端/组件(阿里云 RocketMQ、公司内部 MQ、自研结果存储、自研事件管道)。核心库只提供注册表 + 发现 + 契约校验,不认识任何第三方实现。
状态:v1.0 已落地(D1–D7 全部按建议确认)。实现纪要见 §10;文档引用 docs/design.md §22。
1. 目标 / 非目标¶
目标
- 第三方包
pip install后通过 entry point 自动发现,用户只改一行Config(transport=...); - 私有环境不打包也能用:
app.load_plugins([...])/TASKMQ_PLUGINS=/ CLI--plugins; - 插件只依赖公开契约(Transport 协议 + 数据类型),核心内部结构可以随便重构;
- 插件必须如实声明能力(
supports_leases/supports_workers),框架启动时校验:缺能力要启动即报错或明确降级,不留到运行时诡异失败; - 提供一致性测试套件,第三方后端能自证语义与内建 transport 一致。
非目标
- 不做
sys.path目录扫描、不做动态下载安装(供应链风险,明确拒绝); - 不做 monkey patch 式隐式挂钩(不可读、难排查);
- 不做插件沙箱:插件即代码,信任边界 = 已安装的包 + 显式配置。
2. 现状:插入点在哪(已核对代码)¶
| 位置 | 现状 | 结论 |
|---|---|---|
Config.transport |
已是 Any,校验允许 Transport 实例(hasattr(transport, "enqueue"),config.py:78) |
最简单路径已成立:直接传实例,零基础设施 |
App._make_transport() |
app.py:243,按 scheme 硬编码 memory/sqlite/redis,未知 scheme 抛 ConfigError |
唯一分派口,注册表插这里即可 |
App.add_sink / EventSink |
已协议化,OTel 就是这么接的 | 事件管道已插件友好 |
App.register_codec |
注册自定义类型编解码 | 已有;缺的是自定义序列化器名(_CODEC_FACTORIES 硬编码,protocol.py:469) |
make_pool |
名字硬编码 solo/threads/asyncio/processes(pool.py:235) | 可插件化,但没有真实需求前不做 |
| 子进程 | 进程池用 app_spec="module:attr" 重建 App(execution.py:74) |
插件在子进程也要可加载 → §5 |
| 一致性测试 | 不变量测试散在 tests/test_transport.py 等 |
抽成可复用套件 → §6 |
关键结论:不需要动现有任何一条代码路径,只在 _make_transport() 的「未知 scheme」分支前插入「查注册表」,默认行为完全不变。
3. 设计¶
3.1 一个薄注册表 taskmq/plugins.py¶
TransportFactory = Callable[[TransportOptions], Transport]
CodecFactory = Callable[[CodecRegistry], Codec]
PoolFactory = Callable[[PoolOptions], Pool]
SinkFactory = Callable[[], EventSink]
@dataclasses.dataclass(frozen=True, slots=True)
class TransportOptions:
"""框架交给插件工厂的全部上下文(插件自己解析 URL 里自己的参数)。"""
url: str
codec: Codec # 已按 config.serializer + 类型注册表构造好
codec_registry: CodecRegistry
max_message_bytes: int
idempotency_ttl: float
config: Config # 只读快照
def register_transport(scheme: str, factory: TransportFactory, *, override: bool = False) -> None: ...
def register_codec(name: str, factory: CodecFactory, *, override: bool = False) -> None: ...
def register_sink(name: str, factory: SinkFactory, *, override: bool = False) -> None: ...
def register_pool(name: str, factory: PoolFactory, *, override: bool = False) -> None: ... # 预留
def load_plugins(modules: Iterable[str] | None = None, *, entry_points: bool = True) -> list[str]: ...
def discover(group: str = "taskmq.plugins") -> list[str]: ... # 只返回模块名,不 import
def known_transports() -> dict[str, str]: ... # scheme -> builtin / 插件包名
3.2 解析顺序(App._make_transport)¶
config.transport不是字符串 → 直接返回实例(保留现有语义,给想自己管生命周期的用户);- 内建 scheme(memory / sqlite / redis)→ 走原路径,内建优先;
- 查注册表 → 命中则
factory(TransportOptions(...)); - 仍未命中 → 触发一次 entry point 懒发现(组
taskmq.plugins),再查一次; - 仍无 →
ConfigError,错误信息必须可操作:
未知 transport scheme:rocketmq://
已注册:memory, sqlite, redis(内建)
想接自己的后端?pip install taskmq-rocketmq,或 app.load_plugins(["mycompany.mq"])
细节见 docs/design/plugins.md
懒发现很重要:标准部署(memory/sqlite/redis)不 import 任何第三方包,启动开销不变。
3.3 能力声明与启动校验¶
class RocketMQTransport(Transport):
supports_leases = True # concurrency_key / beat 选主 需要
supports_workers = True # status 的 WORKERS 段需要
- 已有机制:
Worker._validate_capabilities()启动时校验(concurrency_key需要命名租约、async 任务需要 asyncio 池),启动即ConfigError; - 文档写死规则:声明不支持 = 框架不会调用它;插件若谎报(声明 True 但方法抛
NotImplementedError),一致性套件当场抓(§6); - 不支持的项在 CLI 里显式降级显示(例如
supports_workers=False时status打印「(该 transport 不支持 worker 列表)」)。
3.4 URL 约定¶
scheme://host:port/path?k=v 由插件自定义语义,框架只做两件事:原样把 url 交给插件;在文档里列出必须遵守的隔离清单(幂等键、worker 注册表、命名租约都必须落在自己的前缀/命名空间内,避免多环境互踩)。
4. 示例:阿里云 RocketMQ 插件(不改 taskmq 一行)¶
4.1 插件包¶
# taskmq_rocketmq/__init__.py
from urllib.parse import parse_qs, urlparse
from taskmq.plugins import TransportOptions, register_transport
from taskmq.transport import ( # 公开契约:这些名字保证向后兼容
UNSET, DeadLetter, Delivery, JobRecord, JobState, QueueStat, Transport, WorkerInfo,
)
from taskmq.protocol import Envelope, decode_value, encode_value
class RocketMQTransport(Transport):
supports_leases = True # 用 RocketMQ 定时消息实现租约
supports_workers = False # 不维护 worker 注册表,框架会跳过
def __init__(self, *, endpoints: str, topic: str, group: str, codec, registry) -> None:
self._client = _AliyunClient(endpoints, topic, group) # aliyun SDK 只在插件里依赖
self._codec, self._registry = codec, registry
def enqueue(self, env, *, queue=None, delay=0.0, priority=None) -> str:
routed = env.with_routing(queue or env.queue, env.priority if priority is None else priority)
self._client.send(self._codec.encode(routed), delay=delay, key=env.key)
return env.id
# reserve / ack / nack / dead_letter / extend_lease / set_state / get_state /
# queue_stats / reap_expired_leases / reap_expired_jobs / acquire_lease / ...
def _factory(options: TransportOptions) -> Transport:
parsed = urlparse(options.url)
query = parse_qs(parsed.query)
return RocketMQTransport(
endpoints=parsed.netloc or query["endpoints"][0],
topic=parsed.path.lstrip("/") or "taskmq",
group=query.get("group", ["taskmq"])[0],
codec=options.codec,
registry=options.codec_registry,
)
register_transport("rocketmq", _factory)
# taskmq_rocketmq/pyproject.toml —— 被 taskmq 自动发现的唯一要求
[project]
name = "taskmq-rocketmq"
dependencies = ["taskmq>=0.1", "aliyun-sdk-rocketmq>=2"]
[project.entry-points."taskmq.plugins"]
rocketmq = "taskmq_rocketmq"
4.2 用户侧(三种等价写法)¶
# ① 装了包 → 什么都不用做,scheme 自动可用
app = App(Config(transport="rocketmq://rmq.aliyuncs.com:8080/taskmq?group=workers"))
# ② 私有环境不打包:显式加载模块
app = App(Config(transport="rocketmq://..."))
app.load_plugins(["mycompany.mq_adapters.rocketmq"])
# ③ 完全自己管生命周期:直接给实例(现有能力,无需插件机制)
app = App(Config(transport=RocketMQTransport(...)))
export TASKMQ_PLUGINS=mycompany.mq_adapters.rocketmq # 环境变量
taskmq --app myapp:app --plugins mycompany.mq worker -Q email
4.3 自证语义(第三方后端的质量闸门,§6)¶
from taskmq.testing import transport_conformance
def test_rocketmq_matches_taskmq_semantics():
transport_conformance(
lambda: RocketMQTransport(endpoints=..., topic="test", group="t"),
supports={"leases": True, "workers": False},
)
5. 进程边界(不设计好,插件在 worker 子进程里会找不到)¶
- entry points:子进程同样能发现(已安装的包)→ 自动 OK;
- 显式模块:随
ChildTask传下去,子进程load_plugins(payload.plugins)——比只靠环境变量更确定; - 兜底:
TASKMQ_PLUGINS环境变量(文档写明三选一,推荐前两者)。
Worker / 池只需把 app.plugins 透传进 ChildTask(新增字段,有默认值,协议兼容)。
6. 一致性测试套件(本设计里最有价值的部分)¶
把内建 transport 已在测的不变量抽成可复用套件,第三方后端一行调用即可自证:
def transport_conformance(factory, *, supports=..., clock=None, cleanup=None) -> None:
"""跑一遍 Transport 契约;断言失败带场景名,便于定位。"""
| 场景 | 断言 |
|---|---|
| 优先级 | 高优先级先出;同级 FIFO |
| 原子 claim | 8 线程 × 2 连接抢 N 条 → 无重复投递 |
| 幂等 ack | 重复 ack 不报错;迟到 ack 抛 LeaseLost |
| 租约回收 | 过期后重投,deliveries 递增 |
defer |
不消耗 deliveries |
| 让位 | yields 独立递增;到 max_yields 后返回 False |
| 过期 | 过期消息不投递,job 落 EXPIRED |
| 幂等键 | 同 key 第二次不重复入队,返回同一个 job id |
| DLQ | 进 DLQ、可列、可重放(可覆盖优先级) |
| 状态/结果 | set_state/get_state 往返,meta 合并 |
| 命名租约(声明支持时) | 抢/续/放;非持有者续租失败 |
| worker 注册表(声明支持时) | 注册/心跳/注销/列表 |
| job 枚举(声明支持时) | list_jobs 按前缀/状态过滤,新的在前,limit 生效 |
| 能力诚实性 | 声明 supports_*=False 的方法必须抛 TransportError,不能静默返回假数据 |
内建三家 transport 也改成调用同一个套件(测试里一行参数化)——保证「给第三方的契约」就是「我们自己也在跑的契约」,避免文档与实现脱节。
7. 插件 API 的稳定性边界(明确承诺)¶
插件可以依赖(向后兼容;破坏即需要 major/minor + 弃用期)
taskmq.transport:Transport、Delivery、DeadLetter、JobRecord、JobState、QueueStat、UNSET、WorkerInfo、MessageState;taskmq.protocol:Envelope、Codec、CodecRegistry、encode_value/decode_value、encode_payload/decode_payload;taskmq.plugins(本设计新增)、taskmq.errors、taskmq.priority、taskmq.testing;- Transport 的方法签名与语义(§7 契约表)。
插件不得依赖(随时可改):taskmq.worker.* 内部、App._* 私有属性、任何下划线开头的名字、transport 内部表结构。
8. 核心改动清单(约 200 行 + 测试,全部新增)¶
| 文件 | 改动 | 风险 |
|---|---|---|
taskmq/plugins.py |
新增:注册表 + entry point 发现 + TransportOptions | 低(独立模块) |
taskmq/app.py |
未知 scheme 分支前查注册表;App.load_plugins();App.plugins;构造时读 TASKMQ_PLUGINS |
低(内建分支不动) |
taskmq/protocol.py |
_CODEC_FACTORIES 改为「内建 + 可注册」 |
低 |
taskmq/cli.py |
全局 --plugins module,module;错误信息带已注册 scheme |
低 |
worker/execution.py + pool.py |
ChildTask.plugins 并透传 |
低(默认空元组) |
taskmq/testing.py |
transport_conformance()(从现有测试抽取) |
中(要重构现有测试调用它) |
examples/plugin_rocketmq/ |
假实现示例(内存模拟 MQ),证明不改核心也能接 | 低 |
兼容性:默认行为零变化(内建优先、无插件不 import 任何东西);Config(transport=<实例>) 不受影响;ChildTask 新字段有默认值 → 旧调用方不变。
9. ✅ 决策点(已确认)¶
| # | 问题 | 结论 |
|---|---|---|
| D1 | 发现机制 | ✅ 都要:打包走 entry points(组 taskmq.plugins),私有环境走 TASKMQ_PLUGINS / --plugins / app.load_plugins() |
| D2 | 扩展点范围 | ✅ transport + codec + sink 已落地;register_pool 作为预留扩展点(接口在、语义未承诺);result= 独立后端留 Phase 2 |
| D3 | 内建冲突 | ✅ 内建优先;插件覆盖内建必须显式 override=True,否则注册时直接报错 |
| D4 | 子进程传递 | ✅ 写进 ChildTask.plugins(App.plugins 透传),env 作为兜底 |
| D5 | transport_conformance |
✅ 已实现 16 个场景;内建三家 transport 在测试里跑同一套件 |
| D6 | result= |
✅ 本轮不做(结果与状态同源) |
| D7 | limitations 声明 |
✅ 已实现:Transport.limitations + taskmq status 打印 + 一致性套件据此自动跳过对应场景 |
10. 落地纪要(v1.0)¶
代码
| 文件 | 内容 |
|---|---|
taskmq/plugins.py |
注册表(transport/codec/sink/pool 四个扩展点)+ entry point 懒发现 + TransportOptions/PoolOptions + load_plugin(s) + unknown_scheme_message() |
taskmq/app.py |
_make_transport() 查注册表(内建优先、override=True 可覆盖)+ 未知 scheme 懒发现 + App.load_plugins()/App.plugins;include= 与 TASKMQ_PLUGINS= 在 config.validate() 之前加载 |
taskmq/protocol.py |
get_codec 支持插件注册的序列化器名(内建优先,未知时懒发现) |
taskmq/events.py |
build_sink 支持插件注册的 sink 名;错误信息列出全部可选名 |
taskmq/config.py |
pool 允许插件注册名;events 只校验「非空字符串」(名字合法性交给 build_sink → App 构造即失败) |
taskmq/cli.py |
全局 --plugins a,b;status 打印 LIMITATIONS 段(transport 声明的语义降级) |
taskmq/worker/{execution,pool,runner}.py |
ChildTask.plugins + 子进程加载(进程池/可强杀两条路径都覆盖) |
taskmq/transport/base.py |
Transport.limitations(默认空) |
taskmq/testing.py |
transport_conformance():16 个场景(含 capability_honesty / job_listing);假时钟可注入(now=/advance=) |
examples/plugin_rocketmq/ |
示例插件(假 MQ):完整 Transport 实现 + entry point 声明 + README |
测试:tests/test_plugins.py(18 例:注册/内建保护/冲突/未知 scheme 报错/显式加载/env/include/entry point 懒发现与 :attr 钩子/codec/sink/pool/子进程透传/limitations)、
tests/test_conformance.py(8 例:memory + sqlite + redis(Lua 与无 Lua)跑同一套件、声明降级自动跳过、谎报能力被抓、假时钟)、
tests/test_plugin_example.py(7 例:示例插件端到端跑任务、concurrency_key 启动即报错、--plugins 起 worker、未知 scheme 报错可操作)。
实现中确认的两个细节
- sink 名字的合法性必须在
App构造时判,不能放在Config.validate():插件 sink 可能在 validate 之后才注册 (懒发现),所以validate()只保证「非空字符串」,build_sink()负责给出「可选值 + 用 register_sink 注册」的可操作报错 —— 依然是 fail fast。 register()要幂等:模块缓存会让「重复 import」不再产生注册副作用,所以插件应写成if transport_factory("scheme") is None: register_transport(...)(示例插件就是这么写的)。
一致性套件立刻抓到的真问题:示例插件第一版只把消息标成 expired、忘了把 job 置 EXPIRED——
expires 场景当场失败。这正是「契约必须我们自己也在跑」的价值。
示例插件跑套件的结果:12 个场景通过,yield / named_leases / worker_registry /
job_listing 按声明跳过(后者意味着该后端上 DAG 工作流不可用,提交时直接报错)。
6. 文档定稿 + README 增补 + docs/design.md §22 引用。