流式基本流程已通过修改了yield和return的冲突导致的问题

This commit is contained in:
Dong_master
2025-07-04 03:26:44 +08:00
parent 4005a8a3e2
commit 68cdd163d3
8 changed files with 323 additions and 117 deletions

View File

@@ -25,6 +25,8 @@ class MessagePlatformAdapter(metaclass=abc.ABCMeta):
logger: EventLogger
is_stream: bool
def __init__(self, config: dict, ap: app.Application, logger: EventLogger):
"""初始化适配器
@@ -67,6 +69,7 @@ class MessagePlatformAdapter(metaclass=abc.ABCMeta):
message_id: int,
message: platform_message.MessageChain,
quote_origin: bool = False,
is_final: bool = False,
):
"""回复消息(流式输出)
Args:
@@ -114,6 +117,7 @@ class MessagePlatformAdapter(metaclass=abc.ABCMeta):
async def is_stream_output_supported(self) -> bool:
"""是否支持流式输出"""
self.is_stream = False
return False
async def kill(self) -> bool:

View File

@@ -18,6 +18,7 @@ import aiohttp
import lark_oapi.ws.exception
import quart
from lark_oapi.api.im.v1 import *
from lark_oapi.api.cardkit.v1 import *
from .. import adapter
from ...core import app
@@ -348,6 +349,8 @@ class LarkAdapter(adapter.MessagePlatformAdapter):
card_id_dict: dict[str, str]
seq: int
def __init__(self, config: dict, ap: app.Application, logger: EventLogger):
self.config = config
self.ap = ap
@@ -356,6 +359,7 @@ class LarkAdapter(adapter.MessagePlatformAdapter):
self.listeners = {}
self.message_id_to_card_id = {}
self.card_id_dict = {}
self.seq = 0
@self.quart_app.route('/lark/callback', methods=['POST'])
async def lark_callback():
@@ -401,54 +405,79 @@ class LarkAdapter(adapter.MessagePlatformAdapter):
return {'code': 500, 'message': 'error'}
def is_stream_output_supported() -> bool:
async def is_stream_output_supported() -> bool:
is_stream = False
if self.config.get("",None):
if self.config.get("enable-card-reply",None):
is_stream = True
self.is_stream = is_stream
return is_stream
async def create_card_id():
async def create_card_id(message_id):
try:
is_stream = is_stream_output_supported()
is_stream = await is_stream_output_supported()
if is_stream:
self.ap.logger.debug('飞书支持stream输出,创建卡片......')
card_id = ''
if self.card_id_dict:
card_id = [k for k,v in self.card_id_dict.items() if (v+datetime.timedelta(days=14))< datetime.datetime.now()][0]
# card_id = ''
# # if self.card_id_dict:
# # card_id = [k for k,v in self.card_id_dict.items() if (v+datetime.timedelta(days=14))< datetime.datetime.now()][0]
#
# if self.card_id_dict is None:
# # content = {
# # "type": "card_json",
# # "data": {"schema":"2.0","header":{"title":{"content":"bot","tag":"plain_text"}},"body":{"elements":[{"tag":"markdown","content":""}]}}
# # }
# card_data = {"schema":"2.0","header":{"title":{"content":"bot","tag":"plain_text"}},
# "body":{"elements":[{"tag":"markdown","content":""}]},"config": {"streaming_mode": True,
# "streaming_config": {"print_strategy": "fast"}}}
#
# request: CreateCardRequest = CreateCardRequest.builder() \
# .request_body(
# CreateCardRequestBody.builder()
# .type("card_json")
# .data(json.dumps(card_data)) \
# .build()
# ).build()
#
# # 发起请求
# response: CreateCardResponse = self.api_client.cardkit.v1.card.create(request)
#
#
# # 处理失败返回
# if not response.success():
# raise Exception(
# f"client.cardkit.v1.card.create failed, code: {response.code}, msg: {response.msg}, log_id: {response.get_log_id()}, resp: \n{json.dumps(json.loads(response.raw.content), indent=4, ensure_ascii=False)}")
#
# self.ap.logger.debug(f'飞书卡片创建成功,卡片ID: {response.data.card_id}')
# self.card_id_dict[response.data.card_id] = datetime.datetime.now()
#
# card_id = response.data.card_id
card_data = {"schema": "2.0", "header": {"title": {"content": "bot", "tag": "plain_text"}},
"body": {"elements": [{"tag": "markdown", "content": "[思考中.....]","element_id":"markdown_1"}]},
"config": {"streaming_mode": True,
"streaming_config": {"print_strategy": "fast"}}}
if self.card_id_dict is None or card_id == '':
# content = {
# "type": "card_json",
# "data": {"schema":"2.0","header":{"title":{"content":"bot","tag":"plain_text"}},"body":{"elements":[{"tag":"markdown","content":""}]}}
# }
card_data = {"schema":"2.0","header":{"title":{"content":"bot","tag":"plain_text"}},
"body":{"elements":[{"tag":"markdown","content":""}]},"config": {"streaming_mode": True,
"streaming_config": {"print_strategy": "fast"}}}
request: CreateCardRequest = CreateCardRequest.builder() \
.request_body(
CreateCardRequestBody.builder()
.type("card_json")
.data(json.dumps(card_data)) \
.build()
).build()
request: CreateCardRequest = (
CreateCardRequest.builder()
.request_body(
CreateCardRequestBody.builder()
.type("card_json")
.data(json.dumps(card_data))
.build()
)
)
# 发起请求
response: CreateCardResponse = await self.api_client.im.v1.card.create(request)
# 发起请求
response: CreateCardResponse = self.api_client.cardkit.v1.card.create(request)
# 处理失败返回
if not response.success():
raise Exception(
f"client.cardkit.v1.card.create failed, code: {response.code}, msg: {response.msg}, log_id: {response.get_log_id()}, resp: \n{json.dumps(json.loads(response.raw.content), indent=4, ensure_ascii=False)}")
# 处理失败返回
if not response.success():
raise Exception(
f"client.cardkit.v1.card.create failed, code: {response.code}, msg: {response.msg}, log_id: {response.get_log_id()}, resp: \n{json.dumps(json.loads(response.raw.content), indent=4, ensure_ascii=False)}")
self.ap.logger.debug(f'飞书卡片创建成功,卡片ID: {response.data.card_id}')
self.card_id_dict[message_id] = response.data.card_id
self.ap.logger.debug(f'飞书卡片创建成功,卡片ID: {response.data.card_id}')
self.card_id_dict[response.data.card_id] = datetime.datetime.now()
card_id = response.data.card_id
card_id = response.data.card_id
return card_id
except Exception as e:
@@ -458,10 +487,10 @@ class LarkAdapter(adapter.MessagePlatformAdapter):
async def on_message(event: lark_oapi.im.v1.P2ImMessageReceiveV1):
if is_stream_output_supported():
if await is_stream_output_supported():
self.ap.logger.debug('卡片回复模式开启')
# 开启卡片回复模式. 这里可以实现飞书一发消息,马上创建卡片进行回复"思考中..."
card_id = await create_card_id()
card_id = await create_card_id(event.event.message.message_id)
reply_message_id = await self.create_message_card(card_id, event.event.message.message_id)
self.message_id_to_card_id[event.event.message.message_id] = (reply_message_id, time.time())
@@ -500,8 +529,8 @@ class LarkAdapter(adapter.MessagePlatformAdapter):
# TODO 目前只支持卡片模板方式且卡片变量一定是content未来这块要做成可配置
# 发消息马上就会回复显示初始化的content信息即思考中
content = {
'type': 'template',
'data': {'template_id': card_id, 'template_variable': {'content': 'Thinking...'}},
'type': 'card',
'data': {'card_id': card_id, 'template_variable': {'content': 'Thinking...'}},
}
request: ReplyMessageRequest = (
ReplyMessageRequest.builder()
@@ -564,35 +593,49 @@ class LarkAdapter(adapter.MessagePlatformAdapter):
async def reply_message_chunk(
self,
message_source: platform_events.MessageEvent,
message_id: str,
message: platform_message.MessageChain,
quote_origin: bool = False,
is_final: bool = False,
):
"""
回复消息变成更新卡片消息
"""
lark_message = await self.message_converter.yiri2target(message, self.api_client)
if not is_final:
self.seq += 1
text_message = ''
for ele in lark_message[0]:
if ele['tag'] == 'text':
text_message += ele['text']
elif ele['tag'] == 'md':
text_message += ele['text']
print(text_message)
content = {
'type': 'template',
'data': {'template_id': self.config['card_template_id'], 'template_variable': {'content': text_message}},
'type': 'card_json',
'data': {'card_id': self.card_id_dict[message_id], 'elements': {'content': text_message}},
}
request: PatchMessageRequest = (
PatchMessageRequest.builder()
.message_id(self.message_id_to_card_id[message_source.message_chain.message_id][0])
.request_body(PatchMessageRequestBody.builder().content(json.dumps(content)).build())
request: ContentCardElementRequest = ContentCardElementRequest.builder() \
.card_id(self.card_id_dict[message_id]) \
.element_id("markdown_1") \
.request_body(ContentCardElementRequestBody.builder()
# .uuid("a0d69e20-1dd1-458b-k525-dfeca4015204")
.content(text_message)
.sequence(self.seq)
.build()) \
.build()
)
if is_final:
self.seq = 0
# 发起请求
response: PatchMessageResponse = self.api_client.im.v1.message.patch(request)
response: ContentCardElementResponse = self.api_client.cardkit.v1.card_element.content(request)
# 处理失败返回
if not response.success():