From 981bd55f7e15ac9413c72f2ac3909f31a02b6baf Mon Sep 17 00:00:00 2001 From: TangMoonX <1270432252@qq.com> Date: Wed, 8 Jul 2026 08:32:38 +0000 Subject: [PATCH] =?UTF-8?q?feat:=20webhook=E6=96=B9=E5=BC=8F=E5=90=8C?= =?UTF-8?q?=E6=AD=A5issue=E7=8A=B6=E6=80=81=20#1010158081135926589=20#=20R?= =?UTF-8?q?eviewed,=20transaction=20id:=2082978?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- bkmonitor/config/celery/config.py | 6 +- bkmonitor/packages/fta_web/constants.py | 10 ++ bkmonitor/packages/fta_web/issue/resources.py | 61 +++++++- bkmonitor/packages/fta_web/issue/urls.py | 3 +- .../utils/update_issue_status_from_tapd.py | 131 ++++++++++++++++++ bkmonitor/packages/fta_web/tasks.py | 17 +-- 6 files changed, 215 insertions(+), 13 deletions(-) create mode 100644 bkmonitor/packages/fta_web/issue/utils/update_issue_status_from_tapd.py diff --git a/bkmonitor/config/celery/config.py b/bkmonitor/config/celery/config.py index 6aa18e772ed..c955520063d 100644 --- a/bkmonitor/config/celery/config.py +++ b/bkmonitor/config/celery/config.py @@ -146,7 +146,11 @@ class Config: }, "fta_web.tasks.sync_tapd_issue_status": { "task": "fta_web.tasks.sync_tapd_issue_status", - "schedule": crontab(minute="*/10"), + # 实时状态同步以 TAPD webhook (tapd_status_update_callback) 为主路径, + # 本定时任务仅作兜底:补偿 webhook 漏投/回调不可用/历史关联等场景。 + "schedule": crontab( + minute=15, hour="*/2" + ), # 每 2 小时的第 15 分钟执行(00:15, 02:15, 04:15, ...,每天 12 次) "enabled": True, "options": {"queue": "celery_resource"}, }, diff --git a/bkmonitor/packages/fta_web/constants.py b/bkmonitor/packages/fta_web/constants.py index f2164e8672b..8ddf22aac17 100644 --- a/bkmonitor/packages/fta_web/constants.py +++ b/bkmonitor/packages/fta_web/constants.py @@ -766,6 +766,16 @@ def issue_user_oauth(cls) -> str: ) +# TAPD 单据"已完成"状态的关键词 +# 英文关键词(匹配 status_value,即 TAPD 状态的英文 key) +TAPD_COMPLETED_STATUS_EN_KEYWORDS = {"closed", "done", "resolved", "verified"} +# 中文关键词(匹配 display_name,即 TAPD 状态的显示名) +TAPD_COMPLETED_STATUS_CN_KEYWORDS = {"已关闭", "已实现", "已解决", "已完成", "已验证", "done"} + +# TAPD Webhook 支持的单据类型(与 IssueTapdRelation.tapd_type choices 一致) +TAPD_WEBHOOK_SUPPORTED_TYPES = {"story", "bug"} + + class TapdOauthEndpoint: """TAPD OAuth 端点(完整地址,基于 TAPD_OAUTH_BASE_URL)""" diff --git a/bkmonitor/packages/fta_web/issue/resources.py b/bkmonitor/packages/fta_web/issue/resources.py index 894b10b8343..3abba50233b 100644 --- a/bkmonitor/packages/fta_web/issue/resources.py +++ b/bkmonitor/packages/fta_web/issue/resources.py @@ -22,6 +22,7 @@ from django.views.decorators.csrf import csrf_exempt from rest_framework import serializers, exceptions from rest_framework.decorators import api_view +from rest_framework.response import Response from bkm_space.utils import bk_biz_id_to_space_uid from bkmonitor.documents.base import BulkActionType @@ -50,7 +51,7 @@ IssueQueryHandler, ) from fta_web.issue.serializers import IssueSearchSerializer -from fta_web.constants import TapdWorkspaceBindStatus +from fta_web.constants import TapdWorkspaceBindStatus, TAPD_WEBHOOK_SUPPORTED_TYPES from fta_web.issue.utils.tapd import ( save_tapd_token, verify_signed_state, @@ -61,6 +62,7 @@ delete_tapd_token, generate_auth_url, ) +from fta_web.issue.utils.update_issue_status_from_tapd import update_issue_status_from_tapd logger = logging.getLogger("root") @@ -3157,3 +3159,60 @@ def _fail(log_msg, *args): # 5) 302 重定向到 success_url(含 # 的前端地址) return HttpResponseRedirect(success_url) + + +@api_view(["POST"]) +@csrf_exempt +def tapd_status_update_callback(request): + """TAPD Webhook 回调 — 根据tapd的状态实时同步issue状态。 + + TAPD 单据状态变更时推送 update 事件(story::update、bug::update、task::update)。 + 本接口仅处理 story/bug 类型且 status 字段发生变更的事件,其余事件直接忽略。 + + 始终返回 HTTP 200,避免 TAPD 重试。异常仅记录日志,不影响响应。 + """ + try: + data = request.data or {} + + # 1) 解析事件类型 + event = data.get("event", "") + if "::update" not in event: + return Response({"result": True, "message": "ignored non-update event"}, status=200) + + tapd_type = event.split("::")[0].lower() + if tapd_type not in TAPD_WEBHOOK_SUPPORTED_TYPES: + return Response({"result": True, "message": f"ignored unsupported type: {tapd_type}"}, status=200) + + # 2) 检查 status 是否在变更字段中 + change_fields = data.get("change_fields", "") + changed_field_set = {f.strip() for f in change_fields.split(",") if f.strip()} + if "status" not in changed_field_set: + return Response({"result": True, "message": "status not changed, skip"}, status=200) + + # 3) 提取 workspace_id 和 tapd_id + workspace_id = data.get("workspace_id") + tapd_id = str(data.get("id", "")) + + if workspace_id is None or not tapd_id: + logger.warning("TAPD webhook missing workspace_id or id, data=%s", _sanitize_for_log(data)) + return Response({"result": True, "message": "missing workspace_id or tapd id"}, status=200) + + # 4) TAPD 单据状态变更 → 关联 Issue 流转为已解决 + update_issue_status_from_tapd( + workspace_id=workspace_id, + tapd_id=tapd_id, + tapd_type=tapd_type, + ) + + logger.info( + "tapd_status_update_callback processed: event=%s, workspace_id=%s, tapd_id=%s, tapd_type=%s", + event, + workspace_id, + tapd_id, + tapd_type, + ) + return Response({"result": True, "message": "ok"}, status=200) + + except Exception as e: + logger.exception("tapd_status_update_callback error: %s", e) + return Response({"result": False, "message": "internal error"}, status=200) diff --git a/bkmonitor/packages/fta_web/issue/urls.py b/bkmonitor/packages/fta_web/issue/urls.py index 5a2803b8296..cec0e90d291 100644 --- a/bkmonitor/packages/fta_web/issue/urls.py +++ b/bkmonitor/packages/fta_web/issue/urls.py @@ -11,7 +11,7 @@ from django.urls import include, re_path from core.drf_resource.routers import ResourceRouter -from fta_web.issue.resources import tapd_app_install_callback, tapd_user_oauth_callback +from fta_web.issue.resources import tapd_app_install_callback, tapd_user_oauth_callback, tapd_status_update_callback from fta_web.issue.views import IssueViewSet router = ResourceRouter() @@ -20,5 +20,6 @@ urlpatterns = [ re_path(r"^tapd/oauth_callback/$", tapd_user_oauth_callback, name="tapd_user_oauth_callback"), re_path(r"^tapd/app_install_callback/$", tapd_app_install_callback, name="tapd_app_install_callback"), + re_path(r"^tapd/status_update_callback/$", tapd_status_update_callback, name="tapd_status_update_callback"), re_path(r"^", include(router.urls)), ] diff --git a/bkmonitor/packages/fta_web/issue/utils/update_issue_status_from_tapd.py b/bkmonitor/packages/fta_web/issue/utils/update_issue_status_from_tapd.py new file mode 100644 index 00000000000..f611d168073 --- /dev/null +++ b/bkmonitor/packages/fta_web/issue/utils/update_issue_status_from_tapd.py @@ -0,0 +1,131 @@ +""" +TAPD 单据状态变更时,更新关联 Issue 的状态(流转为已解决)。 +""" + +import logging +from concurrent.futures import ThreadPoolExecutor, as_completed + +from bkmonitor.documents.issue import IssueDocument +from bkmonitor.models.issue import IssueTapdRelation +from constants.issue import IssueStatus +from core.drf_resource import api +from fta_web.constants import TAPD_COMPLETED_STATUS_EN_KEYWORDS, TAPD_COMPLETED_STATUS_CN_KEYWORDS + +logger = logging.getLogger("root") + + +def _query_and_check_tapd_status(workspace_id: int, tapd_type: str, tapd_id: str) -> bool: + """查询tapd状态,返回是否为已完成""" + from fta_web.issue.resources import SearchTAPDItemsResource + + items = SearchTAPDItemsResource._query_tapd_items( + tapd_type=tapd_type, + workspace_id=workspace_id, + id=tapd_id, + limit=1, + page=1, + order="created desc", + fields="id,status", + ) + if not items: + return False + item = items[0] + status_value = item.get("status", "") + status_display_name = item.get("status_display_name", "") + if not status_value: + return False + if status_value.lower() in TAPD_COMPLETED_STATUS_EN_KEYWORDS: + return True + if status_display_name and any(keyword in status_display_name for keyword in TAPD_COMPLETED_STATUS_CN_KEYWORDS): + return True + + return False + + +def _filter_resolved(relations: list[dict]) -> list[dict]: + """排除 Issue 已处于"已解决"状态的关联记录。""" + ids = {r["issue_id"] for r in relations} + if not ids: + return relations + + try: + search = ( + IssueDocument.search(all_indices=True) + .filter("terms", _id=list(ids)) + .filter("term", status=IssueStatus.RESOLVED) + .source(["id"]) + ) + resolved_ids = {hit.meta.id for hit in search} + except Exception: + logger.warning("ES query failed, proceeding without pre-filter", exc_info=True) + return relations + + return [r for r in relations if r["issue_id"] not in resolved_ids] + + +def _resolve_one(issue_id: str, bk_biz_id: int) -> dict: + try: + result = api.issue.resolve(bk_biz_id=bk_biz_id, issue_id=issue_id, operator="system") + logger.info("Issue auto-resolved: issue=%s biz=%s", issue_id, bk_biz_id) + return result + except Exception: + logger.warning("Issue resolve failed: issue=%s biz=%s", issue_id, bk_biz_id, exc_info=True) + return {} + + +def _resolve_all(issues: set[tuple[int, str]]) -> tuple[int, int]: + """并发流转 Issue。返回 (resolved, failed)。""" + if not issues: + return 0, 0 + + done = failed = 0 + with ThreadPoolExecutor(max_workers=min(10, len(issues))) as executor: + futures = {executor.submit(_resolve_one, iid, bid): (bid, iid) for bid, iid in issues} + for f in as_completed(futures): + bid, iid = futures[f] + try: + if f.result(): + done += 1 + else: + failed += 1 + except Exception: + logger.warning("resolve panic: biz=%s issue=%s", bid, iid, exc_info=True) + failed += 1 + return done, failed + + +def update_issue_status_from_tapd(workspace_id: int, tapd_id: str, tapd_type: str) -> dict: + """TAPD 单据状态变更为已完成时,将关联 Issue 流转为已解决""" + relations = list( + IssueTapdRelation.objects.filter( + sync_status=True, + workspace_id=workspace_id, + tapd_id=tapd_id, + tapd_type=tapd_type, + ).values("bk_biz_id", "issue_id") + ) + if not relations: + logger.info("no relations, ws=%s tapd=%s", workspace_id, tapd_id) + return {"checked": 0, "resolved": 0, "failed": 0, "skipped": 0} + + try: + completed = _query_and_check_tapd_status(workspace_id, tapd_type, tapd_id) + except Exception: + logger.warning("tapd query failed, ws=%s tapd=%s", workspace_id, tapd_id, exc_info=True) + return {"checked": 0, "resolved": 0, "failed": len(relations), "skipped": 0} + + if not completed: + logger.info("tapd not completed, tapd=%s", tapd_id) + return {"checked": len(relations), "resolved": 0, "failed": 0, "skipped": len(relations)} + + active = _filter_resolved(relations) + skipped = len(relations) - len(active) + if not active: + return {"checked": len(relations), "resolved": 0, "failed": 0, "skipped": len(relations)} + + to_resolve = {(r["bk_biz_id"], r["issue_id"]) for r in active} + done, failed = _resolve_all(to_resolve) + + stats = {"checked": len(relations), "resolved": done, "failed": failed, "skipped": skipped} + logger.info("done: ws=%s tapd=%s stats=%s", workspace_id, tapd_id, stats) + return stats diff --git a/bkmonitor/packages/fta_web/tasks.py b/bkmonitor/packages/fta_web/tasks.py index 9a678fd3c58..5c88d0bcc77 100644 --- a/bkmonitor/packages/fta_web/tasks.py +++ b/bkmonitor/packages/fta_web/tasks.py @@ -21,7 +21,11 @@ from constants.action import ActionPluginType from constants.issue import IssueStatus from core.drf_resource import api, resource -from fta_web.constants import QuickSolutionsConfig +from fta_web.constants import ( + QuickSolutionsConfig, + TAPD_COMPLETED_STATUS_EN_KEYWORDS, + TAPD_COMPLETED_STATUS_CN_KEYWORDS, +) from monitor_web.strategies.user_groups import create_default_notice_group logger = logging.getLogger("celery") @@ -165,13 +169,6 @@ def run_init_builtin_assign_group(bk_biz_id): AlertAssignRule.objects.create(**empty_user_assign_rule) -# TAPD 已完成状态的关键词匹配 -# 英文关键词(匹配 status_value,即 TAPD 状态的英文 key) -_TAPD_COMPLETED_STATUS_EN_KEYWORDS = {"closed", "done", "resolved", "verified"} -# 中文关键词(匹配 display_name,即 TAPD 状态的显示名) -_TAPD_COMPLETED_STATUS_CN_KEYWORDS = {"已关闭", "已实现", "已解决", "已完成", "已验证", "done"} - - def _is_tapd_status_completed(status_value: str, status_display_name: str) -> bool: """判断 TAPD 单据状态是否为已完成状态 @@ -191,11 +188,11 @@ def _is_tapd_status_completed(status_value: str, status_display_name: str) -> bo return False # 1. 先判断英文 key - if status_value.lower() in _TAPD_COMPLETED_STATUS_EN_KEYWORDS: + if status_value.lower() in TAPD_COMPLETED_STATUS_EN_KEYWORDS: return True # 2. 再判断显示名称 display_name - if status_display_name and any(keyword in status_display_name for keyword in _TAPD_COMPLETED_STATUS_CN_KEYWORDS): + if status_display_name and any(keyword in status_display_name for keyword in TAPD_COMPLETED_STATUS_CN_KEYWORDS): return True return False