with one click
mq-middleware
指导 LLM 如何使用 toLink-Rag 项目的 MQ 消息中台进行消息收发、定义新消息类型以及处理多厂商适配逻辑。
Install with Codex or Claude Copy this prompt, paste it into Codex, Claude, or another assistant, and let it review the skill page and install it for you.
Menu
指导 LLM 如何使用 toLink-Rag 项目的 MQ 消息中台进行消息收发、定义新消息类型以及处理多厂商适配逻辑。
Install with Codex or Claude Copy this prompt, paste it into Codex, Claude, or another assistant, and let it review the skill page and install it for you.
Based on SOC occupation classification
契约治理三件套的「值层」。核对同一个物理契约值(MQ topic/group、OSS bucket、消息字段名/别名、内部 HTTP 路径等)在 .env/.env.example/代码生效点/Java 对端多处是否逐字相等,找出配置漂移与死值,防止消息收不到/文件取不到。本 skill 只比对「同一个值在多处是否一致」,不判断结构/语义是否破坏对端(那是结构层,转 contract-guard),也不改文档。
当用户认为当前模块代码实现完毕,且当前分支应为 dev,需要从 dev 基于当前修改创建规范分支、提交并发起合并到 dev 的 GitHub PR 时使用;也用于发布收口,即直接创建 dev -> master 的 release PR,不新建 release 分支。适用于“从 dev 新建分支”“把当前修改提 PR”“实现完成创建 feature/refactor 分支并 PR”“发布新版本”“dev 合并 master”等交付收口场景。本 skill 是交付链终点,并在建分支/提 PR 前执行收口门槛:测试未过、契约文档失同步、acceptance 未提升者拒绝收口。
当用户要求把需求、功能、技术方案、架构改造、故障复盘、项目治理实践或实现过程写成博客/技术文章时必须使用;尤其适用于“写一篇博客”“生成技术博客”“把这个需求写成文章”“根据这个功能写博客”“把项目实现讲清楚”等请求。使用时要基于用户给出的需求和 toLink-Rag 当前仓库的真实代码、文档、契约、配置与测试证据完成分析,默认输出 Markdown 到 `.specs/blog/《博客名称》.md`。文章须采用「少量
把项目里已有的内部组件(如 MQ 中台、解析 pipeline、缓存层、对象存储)抽象成一份「项目自有 skill」,让 AI 每次接入都自动复用该组件的架构边界与约定。读组件真实代码,提炼「架构定位 / 职责边界 / 已落地清单 / 扩展点 / 红线」五要素,按统一原型生成 SKILL.md,登记到 .ai/skills/README.md 注册表并跑校验。
当用户要提 issue、登记 bug、记录新需求时使用;自动识别所属项目,生成结构化 issue 内容,先在 Linear 建主记录、再在 GitHub 建镜像,并双向回链。用户说"提个 issue""记一下这个 bug""把这个需求登记一下""同步到 Linear 和 GitHub""别再依赖 Linear 自动同步"时都应触发,即使没有明确说出"Linear"或"GitHub"。
当用户提出新需求、口头想法、初步框架,或要求"写个 brief / 需求理解 / 业务分析 / 先理清楚需求"时激活;输出面向开发者的 brief.md,覆盖需求摘要、业务流程、涉及模块与影响面(含概念数据模型,但不到物理 schema 与 how)、风险、待确认问题五章;支持开发者审阅后通过对话迭代修订,直到开发者确认冻结,作为后续 acceptance.feature 生成的输入。若用户已基于 brief.md 提出修改、补充、疑问、回答待确认项,继续用本 skill 做迭代收敛。
| name | mq-middleware |
| description | 指导 LLM 如何使用 toLink-Rag 项目的 MQ 消息中台进行消息收发、定义新消息类型以及处理多厂商适配逻辑。 |
| when_to_use | 当用户要求接入 Kafka/RabbitMQ、发送或订阅消息、新增 MQ 消息类型、实现消息消费者或处理多消息队列厂商适配时激活。触发示例:'接入Kafka'、'发送一条消息'、'写个MQ消费者'、'新增消息类型'、'对接RabbitMQ' |
该模块通过 MQFactory 实现多厂商(Kafka/RabbitMQ)切换。LLM 应优先使用 MQService 进行操作,而不是直接实例化 Vendor 适配器。
强制同步要求:
MQFactory、MQService、MQ 配置项。当前项目的职责边界如下:
src/core/mq/message.py:MQ 消息抽象层,定义 MessagePayload、AbstractMessage、统一消息信封序列化和 get_routing_key() 扩展点。src/core/mq/interfaces.py:MQ 厂商能力接口层,定义 IMQSender、IMQReceiver、MQVendorType,业务代码不直接依赖具体 SDK。src/core/mq/exceptions.py:MQ 异常体系,包含连接、发送、消费、配置、序列化异常。src/core/mq/messages/:只放真正的 MQ 业务消息定义,不放 HTTP 请求/响应 DTO。src/core/mq/messages/__init__.py:统一导出当前 MQ 业务消息和 Payload。src/api/schemas/:放 FastAPI 路由使用的请求/响应模型。src/services/mq_service.py:业务侧统一发送/订阅入口,封装 MQFactory -> Sender/Receiver 调用链,支持 send()、send_raw()、subscribe()、start_consuming()、stop_consuming()、close()。src/core/mq/factory.py:注册式单例工厂,根据 MQ_VENDOR 选择 Kafka / RabbitMQ 适配器,缓存 Sender/Receiver,并支持测试场景 reset()。src/core/mq/vendors/kafka/kafka_adapter.py:Kafka 厂商适配器,底层封装 aiokafka,保持 Topic、ConsumerGroup、Offset 语义,消费成功后手动提交 offset。src/core/mq/vendors/kafka/topic_admin.py:Kafka Topic Admin 实现,供 Kafka Topic 管理流程使用。src/core/mq/vendors/rabbitmq_adapter.py:RabbitMQ 厂商适配器,底层封装 aio-pika,保持 Exchange、Queue、Binding、RoutingKey、手动 ACK 语义。src/core/mq/consumers/:消息消费回调实现;当前文档解析消费者位于 src/core/mq/consumers/parse_task_consumer.py,启动入口为 start_parse_consumer()。src/core/mq/topic_admin.py:应用启动阶段可调用的 Kafka Topic Admin 逻辑,当前由 src/main.py 在 MQ_VENDOR=kafka 且 INIT_KAFKA_TOPICS_ON_STARTUP=true 时调用。当前已落地的 MQ 业务消息有 4 类:
src/core/mq/messages/parse_task.py:ParseTaskMessage / ParseTaskPayload,Topic 为 tolink.rag.parse_task,用于文档解析任务投递。src/core/mq/messages/document_delete.py:DocumentDeleteMessage / DocumentDeletePayload,Topic 为 tolink.rag.document_delete,用于清理解析域衍生产物。src/core/mq/messages/token_usage.py:TokenUsageMessage / TokenUsagePayload,Topic 为 tolink.rag.usage_report,用于 LLM 用量上报。src/core/mq/messages/chat_turn.py:ChatTurnMessage / ChatTurnPayload,Topic 为 tolink.rag.chat_turn,用于向 Java 上报对话轮次内容。当前应用启动流程中的 MQ 行为:
src/main.py lifespan 中会初始化 Redis、数据库后进入 MQ 初始化逻辑。settings.MQ_VENDOR.lower() == "kafka" 且 settings.INIT_KAFKA_TOPICS_ON_STARTUP 为 true 时,调用 src/core/mq/topic_admin.py::ensure_topics()。parse_task 与 document_delete 两个消费者,然后统一启动 MQService 消费。不要把消息模型拆成 payload.py / message.py 两个文件,也不要把 HTTP DTO 放进 src/core/mq/messages/。
当用户要求“发送某某通知”或“触发某项异步任务”时:
src/core/mq/messages/ 下是否已有对应的消息模型。MQService().send(YourMessage.build(...))。src/core/mq/messages/your_event.py。当用户要求“监听消息”或“处理 MQ 任务”时:
MQService().subscribe(topic, group_id, callback)。callback 是一个 async 函数。MQService().start_consuming() 才会开始拉取消息。src/core/mq/consumers/。src/core/mq/consumers/parse_task_consumer.py。新增 MQ 消息时,按“一个业务消息一个文件”组织。例如:
src/core/mq/messages/
parse_task.py
document_delete.py
token_usage.py
chat_turn.py
your_event.py
每个文件内部同时定义:
YourPayloadYourMessageMQReceiver Protocol不要新增以下结构:
your_payload.pyyour_message.pyhttp_models.py如果需要新增业务消息,请按以下结构生成代码:
from src.core.mq.message import AbstractMessage, MessagePayload
from pydantic import Field
from typing import Protocol
class YourPayload(MessagePayload):
# 定义具体字段
biz_id: str = Field(..., title="业务ID")
class YourMessage(AbstractMessage):
MQ_NAME = "your.topic.name"
MQ_TYPE = "YOUR_TYPE"
def __init__(self, payload: YourPayload):
self._payload = payload
@classmethod
def get_mq_name(cls): return cls.MQ_NAME
@classmethod
def get_mq_type(cls): return cls.MQ_TYPE
def get_payload(self): return self._payload
@classmethod
def build(cls, **kwargs):
return cls(payload=YourPayload(**kwargs))
@classmethod
def parse_msg(cls, raw: str) -> YourPayload:
envelope = cls.deserialize_envelope(raw)
return YourPayload(**envelope["payload"])
class MQReceiver(Protocol):
async def on_your_event(self, payload: "YourPayload") -> None: ...
补充要求:
MQ_NAME 使用稳定 topic 名称,例如 tolink.rag.your_eventMQ_TYPE 使用稳定枚举式字符串,例如 YOUR_EVENTget_routing_key()当前项目支持在应用启动时可选初始化 Kafka Topics:
src/main.py 的 lifespan 中调用 src/core/mq/topic_admin.py相关约定:
.env / settings 中的 INIT_KAFKA_TOPICS_ON_STARTUPMQ_VENDOR=kafka 且开关为 true 时才执行初始化当前默认 Topic:
tolink.rag.parse_tasktolink.rag.usage_reporttolink.rag.chat_turntolink.rag.document_delete.env 的 MQ_VENDOR 字段。aiokafkaconfluent-kafkaaio-pikaKAFKA_BOOTSTRAP_SERVERSKAFKA_SECURITY_PROTOCOLKAFKA_SASL_MECHANISMKAFKA_SASL_USERNAMEKAFKA_SASL_PASSWORDKAFKA_MAX_POLL_INTERVAL_MSINIT_KAFKA_TOPICS_ON_STARTUPRABBITMQ_URLRABBITMQ_EXCHANGE_NAMERABBITMQ_EXCHANGE_TYPERABBITMQ_PREFETCH_COUNT调试时优先检查:
MQService logsrc/core/mq/vendors/kafka/kafka_adapter.pysrc/core/mq/vendors/rabbitmq_adapter.pysrc/core/mq/topic_admin.pysrc/core/mq/vendors/kafka/topic_admin.py