forked from erp-dev/erp
feat: message-api for mission via wecomm agent
This commit is contained in:
@@ -5,8 +5,8 @@ from mission.models import Mission, MissionCategory, MissionParticipant, Mission
|
||||
|
||||
@admin.register(MissionCategory)
|
||||
class MissionCategoryAdmin(admin.ModelAdmin):
|
||||
list_display = ["id", "merchant", "name", "created_at"]
|
||||
list_filter = ["merchant", "created_at"]
|
||||
list_display = ["id", "merchant", "name", "payload_processor", "created_at"]
|
||||
list_filter = ["merchant", "payload_processor", "created_at"]
|
||||
search_fields = ["name"]
|
||||
readonly_fields = ["created_at", "updated_at"]
|
||||
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
import logging
|
||||
|
||||
from mission.payload_processors import apply_mission_payload_processor
|
||||
from notifier.models import NotificationEventKeyEnum
|
||||
from notifier.services import enqueue_notification_event
|
||||
|
||||
@@ -26,6 +27,7 @@ def _build_mission_payload(mission) -> dict:
|
||||
"participant_names_display": "、".join(participant_names) if participant_names else "无",
|
||||
"content_type": content_type or "",
|
||||
"content_id": mission.content_id or "",
|
||||
"payload_processor": getattr(mission.category, "payload_processor", "") or "",
|
||||
}
|
||||
|
||||
|
||||
@@ -58,11 +60,15 @@ def _enqueue(*, event_key: str, merchant_id: int, payload: dict) -> None:
|
||||
|
||||
|
||||
def on_mission_created(sender, instance, created_by=None, **kwargs):
|
||||
payload = {
|
||||
**_build_mission_payload(instance),
|
||||
"created_by_id": getattr(created_by, "id", None),
|
||||
"created_by_name": getattr(created_by, "name", ""),
|
||||
}
|
||||
payload = apply_mission_payload_processor(
|
||||
mission=instance,
|
||||
event_key=NotificationEventKeyEnum.MISSION_CREATED,
|
||||
payload={
|
||||
**_build_mission_payload(instance),
|
||||
"created_by_id": getattr(created_by, "id", None),
|
||||
"created_by_name": getattr(created_by, "name", ""),
|
||||
},
|
||||
)
|
||||
_enqueue(
|
||||
event_key=NotificationEventKeyEnum.MISSION_CREATED,
|
||||
merchant_id=instance.merchant_id,
|
||||
@@ -72,12 +78,16 @@ def on_mission_created(sender, instance, created_by=None, **kwargs):
|
||||
|
||||
def on_mission_replied(sender, instance, mission=None, responder=None, **kwargs):
|
||||
mission = mission or instance.mission
|
||||
payload = {
|
||||
**_build_mission_payload(mission),
|
||||
**_build_reply_payload(instance),
|
||||
"responder_id": getattr(responder, "id", None),
|
||||
"responder_name": getattr(responder, "name", ""),
|
||||
}
|
||||
payload = apply_mission_payload_processor(
|
||||
mission=mission,
|
||||
event_key=NotificationEventKeyEnum.MISSION_REPLIED,
|
||||
payload={
|
||||
**_build_mission_payload(mission),
|
||||
**_build_reply_payload(instance),
|
||||
"responder_id": getattr(responder, "id", None),
|
||||
"responder_name": getattr(responder, "name", ""),
|
||||
},
|
||||
)
|
||||
_enqueue(
|
||||
event_key=NotificationEventKeyEnum.MISSION_REPLIED,
|
||||
merchant_id=mission.merchant_id,
|
||||
@@ -93,6 +103,11 @@ def on_mission_completed(sender, instance, completed_by=None, reply=None, **kwar
|
||||
}
|
||||
if reply is not None:
|
||||
payload.update(_build_reply_payload(reply))
|
||||
payload = apply_mission_payload_processor(
|
||||
mission=instance,
|
||||
event_key=NotificationEventKeyEnum.MISSION_COMPLETED,
|
||||
payload=payload,
|
||||
)
|
||||
_enqueue(
|
||||
event_key=NotificationEventKeyEnum.MISSION_COMPLETED,
|
||||
merchant_id=instance.merchant_id,
|
||||
@@ -102,13 +117,17 @@ def on_mission_completed(sender, instance, completed_by=None, reply=None, **kwar
|
||||
|
||||
def on_mission_reply_rejected(sender, instance, mission=None, rejected_by=None, reason=None, **kwargs):
|
||||
mission = mission or instance.mission
|
||||
payload = {
|
||||
**_build_mission_payload(mission),
|
||||
**_build_reply_payload(instance),
|
||||
"rejected_by_id": getattr(rejected_by, "id", None),
|
||||
"rejected_by_name": getattr(rejected_by, "name", ""),
|
||||
"reason": reason or "",
|
||||
}
|
||||
payload = apply_mission_payload_processor(
|
||||
mission=mission,
|
||||
event_key=NotificationEventKeyEnum.MISSION_REPLY_REJECTED,
|
||||
payload={
|
||||
**_build_mission_payload(mission),
|
||||
**_build_reply_payload(instance),
|
||||
"rejected_by_id": getattr(rejected_by, "id", None),
|
||||
"rejected_by_name": getattr(rejected_by, "name", ""),
|
||||
"reason": reason or "",
|
||||
},
|
||||
)
|
||||
_enqueue(
|
||||
event_key=NotificationEventKeyEnum.MISSION_REPLY_REJECTED,
|
||||
merchant_id=mission.merchant_id,
|
||||
@@ -117,13 +136,17 @@ def on_mission_reply_rejected(sender, instance, mission=None, rejected_by=None,
|
||||
|
||||
|
||||
def on_mission_reopened(sender, instance, reopened_by=None, rejected_reply_ids=None, **kwargs):
|
||||
payload = {
|
||||
**_build_mission_payload(instance),
|
||||
"reopened_by_id": getattr(reopened_by, "id", None),
|
||||
"reopened_by_name": getattr(reopened_by, "name", ""),
|
||||
"rejected_reply_ids": rejected_reply_ids or [],
|
||||
"rejected_reply_ids_display": ", ".join(str(reply_id) for reply_id in (rejected_reply_ids or [])) or "无",
|
||||
}
|
||||
payload = apply_mission_payload_processor(
|
||||
mission=instance,
|
||||
event_key=NotificationEventKeyEnum.MISSION_REOPENED,
|
||||
payload={
|
||||
**_build_mission_payload(instance),
|
||||
"reopened_by_id": getattr(reopened_by, "id", None),
|
||||
"reopened_by_name": getattr(reopened_by, "name", ""),
|
||||
"rejected_reply_ids": rejected_reply_ids or [],
|
||||
"rejected_reply_ids_display": ", ".join(str(reply_id) for reply_id in (rejected_reply_ids or [])) or "无",
|
||||
},
|
||||
)
|
||||
_enqueue(
|
||||
event_key=NotificationEventKeyEnum.MISSION_REOPENED,
|
||||
merchant_id=instance.merchant_id,
|
||||
@@ -132,12 +155,16 @@ def on_mission_reopened(sender, instance, reopened_by=None, rejected_reply_ids=N
|
||||
|
||||
|
||||
def on_mission_cancelled(sender, instance, cancelled_by=None, **kwargs):
|
||||
payload = {
|
||||
**_build_mission_payload(instance),
|
||||
"cancelled_by_id": getattr(cancelled_by, "id", None),
|
||||
"cancelled_by_name": getattr(cancelled_by, "name", ""),
|
||||
"cancelled_at": instance.cancelled_at.isoformat() if instance.cancelled_at else "",
|
||||
}
|
||||
payload = apply_mission_payload_processor(
|
||||
mission=instance,
|
||||
event_key=NotificationEventKeyEnum.MISSION_CANCELLED,
|
||||
payload={
|
||||
**_build_mission_payload(instance),
|
||||
"cancelled_by_id": getattr(cancelled_by, "id", None),
|
||||
"cancelled_by_name": getattr(cancelled_by, "name", ""),
|
||||
"cancelled_at": instance.cancelled_at.isoformat() if instance.cancelled_at else "",
|
||||
},
|
||||
)
|
||||
_enqueue(
|
||||
event_key=NotificationEventKeyEnum.MISSION_CANCELLED,
|
||||
merchant_id=instance.merchant_id,
|
||||
|
||||
22
mission/migrations/0007_missioncategory_payload_processor.py
Normal file
22
mission/migrations/0007_missioncategory_payload_processor.py
Normal file
@@ -0,0 +1,22 @@
|
||||
from django.db import migrations, models
|
||||
|
||||
|
||||
class Migration(migrations.Migration):
|
||||
|
||||
dependencies = [
|
||||
("mission", "0006_mission_and_reply_extra"),
|
||||
]
|
||||
|
||||
operations = [
|
||||
migrations.AddField(
|
||||
model_name="missioncategory",
|
||||
name="payload_processor",
|
||||
field=models.CharField(
|
||||
blank=True,
|
||||
choices=[("structured_description_v1", "结构化描述增强(v1)")],
|
||||
default="",
|
||||
max_length=50,
|
||||
verbose_name="payload 增强器",
|
||||
),
|
||||
),
|
||||
]
|
||||
@@ -6,6 +6,10 @@ from django.utils import timezone
|
||||
from flower.common import ModelBase
|
||||
|
||||
|
||||
class MissionPayloadProcessorEnum(models.TextChoices):
|
||||
STRUCTURED_DESCRIPTION_V1 = "structured_description_v1", "结构化描述增强(v1)"
|
||||
|
||||
|
||||
class MissionCategory(ModelBase):
|
||||
"""任务分类字典,按商户隔离。"""
|
||||
|
||||
@@ -17,6 +21,13 @@ class MissionCategory(ModelBase):
|
||||
verbose_name="所属商户",
|
||||
)
|
||||
name = models.CharField(max_length=50, verbose_name="分类名称")
|
||||
payload_processor = models.CharField(
|
||||
max_length=50,
|
||||
blank=True,
|
||||
default="",
|
||||
choices=MissionPayloadProcessorEnum.choices,
|
||||
verbose_name="payload 增强器",
|
||||
)
|
||||
|
||||
def __str__(self):
|
||||
return self.name
|
||||
|
||||
86
mission/payload_processors.py
Normal file
86
mission/payload_processors.py
Normal file
@@ -0,0 +1,86 @@
|
||||
import re
|
||||
|
||||
from mission.models import MissionPayloadProcessorEnum
|
||||
|
||||
|
||||
STRUCTURED_DESCRIPTION_IMAGE_PREFIX = "款式图:"
|
||||
|
||||
|
||||
def _extract_labeled_image_url(lines: list[str]) -> tuple[list[str], str]:
|
||||
kept_lines = []
|
||||
image_url = ""
|
||||
for line in lines:
|
||||
stripped_line = line.strip()
|
||||
if stripped_line.startswith(STRUCTURED_DESCRIPTION_IMAGE_PREFIX):
|
||||
candidate = stripped_line.removeprefix(STRUCTURED_DESCRIPTION_IMAGE_PREFIX).strip()
|
||||
if candidate:
|
||||
image_url = candidate
|
||||
continue
|
||||
kept_lines.append(line)
|
||||
return kept_lines, image_url
|
||||
|
||||
|
||||
def _split_structured_description(description: str) -> tuple[str, str, str, str]:
|
||||
lines = [line.strip() for line in (description or "").splitlines()]
|
||||
if not lines:
|
||||
return "", "", "", ""
|
||||
|
||||
body_lines = lines[1:] if len(lines) > 1 else []
|
||||
while body_lines and not body_lines[0]:
|
||||
body_lines.pop(0)
|
||||
|
||||
body_lines, image_url = _extract_labeled_image_url(body_lines)
|
||||
|
||||
url = ""
|
||||
if body_lines:
|
||||
last_line = body_lines[-1]
|
||||
match = re.search(r"https?://\S+", last_line)
|
||||
if match is not None:
|
||||
url = match.group(0)
|
||||
body_lines = body_lines[:-1]
|
||||
|
||||
while body_lines and not body_lines[-1]:
|
||||
body_lines.pop()
|
||||
|
||||
split_index = next((index for index, line in enumerate(body_lines) if not line), None)
|
||||
if split_index is None:
|
||||
title_lines = body_lines
|
||||
description_lines = []
|
||||
else:
|
||||
title_lines = body_lines[:split_index]
|
||||
description_lines = body_lines[split_index + 1 :]
|
||||
while description_lines and not description_lines[0]:
|
||||
description_lines.pop(0)
|
||||
|
||||
title = "\n".join(title_lines).strip()
|
||||
parsed_description = "\n".join(description_lines).strip()
|
||||
return title, parsed_description, url, image_url
|
||||
|
||||
|
||||
def enhance_structured_description_v1(*, mission, payload: dict, event_key: str) -> dict:
|
||||
title, parsed_description, url, image_url = _split_structured_description(mission.description)
|
||||
return {
|
||||
**payload,
|
||||
"parsed_description_title": title,
|
||||
"parsed_description_body": parsed_description,
|
||||
"parsed_description_url": url,
|
||||
"parsed_description_image_url": image_url,
|
||||
}
|
||||
|
||||
|
||||
PAYLOAD_PROCESSORS = {
|
||||
MissionPayloadProcessorEnum.STRUCTURED_DESCRIPTION_V1: enhance_structured_description_v1,
|
||||
}
|
||||
|
||||
|
||||
def apply_mission_payload_processor(*, mission, payload: dict, event_key: str) -> dict:
|
||||
processor_key = getattr(mission.category, "payload_processor", "") or ""
|
||||
if not processor_key:
|
||||
return payload
|
||||
|
||||
processor = PAYLOAD_PROCESSORS.get(processor_key)
|
||||
if processor is None:
|
||||
return payload
|
||||
|
||||
enhanced_payload = processor(mission=mission, payload=payload, event_key=event_key)
|
||||
return enhanced_payload if isinstance(enhanced_payload, dict) else payload
|
||||
110
mission/tests.py
110
mission/tests.py
@@ -1,5 +1,6 @@
|
||||
from django.test import TestCase
|
||||
from django.contrib.contenttypes.models import ContentType
|
||||
from types import SimpleNamespace
|
||||
from unittest.mock import patch
|
||||
from django.utils import timezone
|
||||
|
||||
@@ -23,6 +24,9 @@ from mission.signals import (
|
||||
mission_replied,
|
||||
mission_reply_rejected,
|
||||
)
|
||||
from mission.payload_processors import _split_structured_description
|
||||
from notifier.models import NotificationEventKeyEnum, Notifier, NotifierChannelEnum, NotifierRoute
|
||||
from notifier.services import dispatch_notification_event
|
||||
|
||||
|
||||
class MissionModelTestCase(TestCase):
|
||||
@@ -546,6 +550,26 @@ class MissionModelTestCase(TestCase):
|
||||
|
||||
self.assertEqual(received, [(Mission, self.mission.id, self.creator.id)])
|
||||
|
||||
def test_split_structured_description_v1(self):
|
||||
title, body, url, image_url = _split_structured_description(
|
||||
"系统首行\n标题一\n款式图:https://images.yuwen.cloud/abc.jpg\n标题二\n\n正文一\n正文二\n手机端链接: https://example.com/detail"
|
||||
)
|
||||
|
||||
self.assertEqual(title, "标题一\n标题二")
|
||||
self.assertEqual(body, "正文一\n正文二")
|
||||
self.assertEqual(url, "https://example.com/detail")
|
||||
self.assertEqual(image_url, "https://images.yuwen.cloud/abc.jpg")
|
||||
|
||||
def test_split_structured_description_v1_without_blank_line_uses_title_only(self):
|
||||
title, body, url, image_url = _split_structured_description(
|
||||
"系统首行\n标题一\n款式图:https://images.yuwen.cloud/abc.jpg\n标题二\n手机端链接: https://example.com/detail"
|
||||
)
|
||||
|
||||
self.assertEqual(title, "标题一\n标题二")
|
||||
self.assertEqual(body, "")
|
||||
self.assertEqual(url, "https://example.com/detail")
|
||||
self.assertEqual(image_url, "https://images.yuwen.cloud/abc.jpg")
|
||||
|
||||
@patch("mission.handlers.enqueue_notification_event")
|
||||
def test_mission_created_handler_enqueues_notifier_event(self, mock_enqueue):
|
||||
with self.captureOnCommitCallbacks(execute=True):
|
||||
@@ -556,6 +580,28 @@ class MissionModelTestCase(TestCase):
|
||||
self.assertEqual(mock_enqueue.call_args.kwargs["merchant_id"], self.merchant.id)
|
||||
self.assertEqual(mock_enqueue.call_args.kwargs["payload"]["mission_id"], mission.id)
|
||||
|
||||
@patch("mission.handlers.enqueue_notification_event")
|
||||
def test_mission_created_handler_applies_category_payload_processor(self, mock_enqueue):
|
||||
self.default_category.payload_processor = "structured_description_v1"
|
||||
self.default_category.save(update_fields=["payload_processor", "updated_at"])
|
||||
|
||||
with self.captureOnCommitCallbacks(execute=True):
|
||||
mission = create_mission(
|
||||
creator=self.creator,
|
||||
description=(
|
||||
"系统首行\n标题一\n款式图:https://images.yuwen.cloud/abc.jpg\n标题二\n\n正文一\n正文二\n"
|
||||
"手机端链接: https://example.com/detail"
|
||||
),
|
||||
)
|
||||
|
||||
payload = mock_enqueue.call_args.kwargs["payload"]
|
||||
self.assertEqual(payload["mission_id"], mission.id)
|
||||
self.assertEqual(payload["payload_processor"], "structured_description_v1")
|
||||
self.assertEqual(payload["parsed_description_title"], "标题一\n标题二")
|
||||
self.assertEqual(payload["parsed_description_body"], "正文一\n正文二")
|
||||
self.assertEqual(payload["parsed_description_url"], "https://example.com/detail")
|
||||
self.assertEqual(payload["parsed_description_image_url"], "https://images.yuwen.cloud/abc.jpg")
|
||||
|
||||
@patch("mission.handlers.enqueue_notification_event")
|
||||
def test_create_ending_reply_handler_enqueues_replied_and_completed_notifications(self, mock_enqueue):
|
||||
with self.captureOnCommitCallbacks(execute=True):
|
||||
@@ -614,3 +660,67 @@ class MissionModelTestCase(TestCase):
|
||||
mock_enqueue.assert_called_once()
|
||||
self.assertEqual(mock_enqueue.call_args.kwargs["event_key"], "mission.cancelled")
|
||||
self.assertEqual(mock_enqueue.call_args.kwargs["payload"]["mission_id"], self.mission.id)
|
||||
|
||||
@patch("notifier.backends.send_message_api_news_to_agents")
|
||||
@patch("notifier.tasks.dispatch_notification_event_task.delay")
|
||||
def test_mission_created_can_flow_to_message_api_news_with_structured_description_template(
|
||||
self,
|
||||
mock_delay,
|
||||
mock_send_news,
|
||||
):
|
||||
self.default_category.payload_processor = "structured_description_v1"
|
||||
self.default_category.save(update_fields=["payload_processor", "updated_at"])
|
||||
notifier = Notifier.objects.create(
|
||||
merchant=self.merchant,
|
||||
name="结构化描述图文通知",
|
||||
channel=NotifierChannelEnum.MESSAGE_API,
|
||||
template_key="mission_structured_description_news",
|
||||
config={
|
||||
"agent_ids": [1000007],
|
||||
"image_url": "https://cdn.example.com/covers/mission-news.png",
|
||||
},
|
||||
)
|
||||
NotifierRoute.objects.create(
|
||||
merchant=self.merchant,
|
||||
notifier=notifier,
|
||||
event_key=NotificationEventKeyEnum.MISSION_CREATED,
|
||||
mission_category=self.default_category,
|
||||
)
|
||||
mock_send_news.return_value = [
|
||||
{
|
||||
"agent_id": 1000007,
|
||||
"ok": True,
|
||||
"response": {"errcode": 0, "errmsg": "ok"},
|
||||
"error": None,
|
||||
}
|
||||
]
|
||||
|
||||
def inline_delay(*, event_key, merchant_id, payload):
|
||||
dispatch_notification_event(
|
||||
event_key=event_key,
|
||||
merchant_id=merchant_id,
|
||||
payload=payload,
|
||||
)
|
||||
return SimpleNamespace(id="task-inline-1")
|
||||
|
||||
mock_delay.side_effect = inline_delay
|
||||
|
||||
with self.captureOnCommitCallbacks(execute=True):
|
||||
create_mission(
|
||||
creator=self.creator,
|
||||
category=self.default_category,
|
||||
description=(
|
||||
"系统首行\n标题一\n款式图:https://images.yuwen.cloud/abc.jpg\n标题二\n\n正文一\n正文二\n"
|
||||
"手机端链接: https://example.com/detail"
|
||||
),
|
||||
)
|
||||
|
||||
mock_send_news.assert_called_once()
|
||||
self.assertEqual(mock_send_news.call_args.kwargs["agent_ids"], [1000007])
|
||||
self.assertEqual(mock_send_news.call_args.kwargs["title"], "标题一\n标题二")
|
||||
self.assertEqual(mock_send_news.call_args.kwargs["description"], "正文一\n正文二")
|
||||
self.assertEqual(mock_send_news.call_args.kwargs["url"], "https://example.com/detail")
|
||||
self.assertEqual(
|
||||
mock_send_news.call_args.kwargs["image_url"],
|
||||
"https://images.yuwen.cloud/abc.jpg",
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user