ワンクリックで
mq-middleware
指导 LLM 如何使用 toLink-Rag 项目的 MQ 消息中台进行消息收发、定义新消息类型以及处理多厂商适配逻辑。
Codex または Claude でインストール この Prompt をコピーして Codex、Claude、または他のアシスタントに貼り付けると、Skill ページを確認してインストールできます。
メニュー
指导 LLM 如何使用 toLink-Rag 项目的 MQ 消息中台进行消息收发、定义新消息类型以及处理多厂商适配逻辑。
Codex または Claude でインストール この Prompt をコピーして Codex、Claude、または他のアシスタントに貼り付けると、Skill ページを確認してインストールできます。
契约治理三件套的「值层」。核对同一个物理契约值(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 做迭代收敛。
SOC 職業分類に基づく
| 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