Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 5 additions & 1 deletion bkmonitor/config/celery/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"},
},
Expand Down
10 changes: 10 additions & 0 deletions bkmonitor/packages/fta_web/constants.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)"""

Expand Down
61 changes: 60 additions & 1 deletion bkmonitor/packages/fta_web/issue/resources.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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,
Expand All @@ -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")

Expand Down Expand Up @@ -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)
3 changes: 2 additions & 1 deletion bkmonitor/packages/fta_web/issue/urls.py
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand All @@ -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)),
]
Original file line number Diff line number Diff line change
@@ -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
17 changes: 7 additions & 10 deletions bkmonitor/packages/fta_web/tasks.py
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down Expand Up @@ -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 单据状态是否为已完成状态

Expand All @@ -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
Expand Down