From 3b0969cb09f92d0dfff95c901fdb8b7f61b5fc56 Mon Sep 17 00:00:00 2001 From: Wenruli Date: Wed, 5 Aug 2026 20:45:49 +0800 Subject: [PATCH] feat: update dashboard data pipeline --- docker/bisheng/config/config.yaml | 10 +- docker/bisheng/config/config_dev.yaml | 10 +- docs/dashboard-dataset-metric-calculation.md | 130 +++++++ features/v2.5.0-sg/release-contract.md | 3 + .../bisheng/api/services/knowledge_imp.py | 2 +- src/backend/bisheng/config_3002.yaml | 10 +- src/backend/bisheng/config_3003.yaml | 10 +- .../bisheng/core/config/celery_queues.py | 61 +++ src/backend/bisheng/core/config/settings.py | 20 +- .../domain/services/favorite_notify.py | 2 +- .../services/knowledge_migration_service.py | 4 +- .../services/knowledge_space_service.py | 2 +- .../portal_hot_search_admin_service.py | 6 +- .../org_sync/api/endpoints/sync_exec.py | 2 +- .../primary_department_change_coordinator.py | 2 +- .../telemetry_search/domain/init_dataset.py | 22 +- .../domain/services/dashboard.py | 139 ++++++- src/backend/bisheng/worker/config.py | 9 +- .../worker/knowledge/document_projection.py | 14 +- .../worker/knowledge/portal_hot_search.py | 4 +- .../worker/knowledge/portal_recommendation.py | 20 +- .../worker/org_sync/reconcile_tasks.py | 2 +- src/backend/bisheng/worker/org_sync/tasks.py | 4 +- .../permission/department_transfer_cleanup.py | 2 +- src/backend/celerybeat-schedule.db | Bin 7737344 -> 7737344 bytes ...reconcile_knowledge_document_projection.py | 4 +- .../test_knowledge_parse_queue_routing.py | 173 +++++++++ .../test_favorite_notification_event.py | 2 +- ...ge_document_permission_reconcile_worker.py | 2 +- .../test_portal_hot_search_worker.py | 2 +- .../test_portal_recommendation_worker.py | 2 +- ...test_department_transfer_cleanup_worker.py | 4 +- .../test_cleanup_service.py | 2 +- .../test_dashboard_enum_labels.py | 348 ++++++++++++++++++ src/backend/test/test_realtime_dashboard.py | 175 +++++++++ .../Dashboard/components/charts/BaseChart.tsx | 22 +- .../platform/src/test/pieChartTooltip.test.ts | 66 ++++ 37 files changed, 1207 insertions(+), 85 deletions(-) create mode 100644 docs/dashboard-dataset-metric-calculation.md create mode 100644 src/backend/bisheng/core/config/celery_queues.py create mode 100644 src/backend/test/celery/test_knowledge_parse_queue_routing.py create mode 100644 src/backend/test/telemetry_search/test_dashboard_enum_labels.py create mode 100644 src/frontend/platform/src/test/pieChartTooltip.test.ts diff --git a/docker/bisheng/config/config.yaml b/docker/bisheng/config/config.yaml index 5b976b0ba..d48df82bc 100644 --- a/docker/bisheng/config/config.yaml +++ b/docker/bisheng/config/config.yaml @@ -66,10 +66,16 @@ celery_redis_url: "redis://redis:6379/2" celery_task: # 对celery熟悉的用户可以自定义配置任务的路由,启动不同类型的worker处理不同类型的异步任务。注意工作流的执行只能在一个进程内!!! task_routers: + bisheng.worker.knowledge.file_title_worker.extract_knowledge_file_title_celery: + queue: knowledge_celery + bisheng.worker.knowledge.file_worker.parse_knowledge_file_celery: + queue: knowledge_celery + bisheng.worker.knowledge.file_worker.retry_knowledge_file_celery: + queue: knowledge_celery bisheng.worker.knowledge.pdf_artifact_worker.generate_knowledge_file_pdf_celery: queue: knowledge_pdf_celery - bisheng.worker.knowledge.*: # 知识库文件处理相关任务 - queue: knowledge_celery + bisheng.worker.knowledge.*: # 其他知识库后台任务 + queue: celery bisheng.worker.workflow.*: # 工作流相关任务 queue: workflow_celery diff --git a/docker/bisheng/config/config_dev.yaml b/docker/bisheng/config/config_dev.yaml index c15e6863a..aa7a4a287 100644 --- a/docker/bisheng/config/config_dev.yaml +++ b/docker/bisheng/config/config_dev.yaml @@ -40,8 +40,16 @@ celery_task: knowledge_file_time_limit: null # 每秒处理10个文件 null: 表示无限制 # 对celery熟悉的用户可以自定义配置任务的路由,启动不同类型的worker处理不同类型的异步任务,注意工作流的执行只能在一个进程内 task_routers: - bisheng.worker.knowledge.*: # 知识库文件处理相关任务 + bisheng.worker.knowledge.file_title_worker.extract_knowledge_file_title_celery: queue: knowledge_celery + bisheng.worker.knowledge.file_worker.parse_knowledge_file_celery: + queue: knowledge_celery + bisheng.worker.knowledge.file_worker.retry_knowledge_file_celery: + queue: knowledge_celery + bisheng.worker.knowledge.pdf_artifact_worker.generate_knowledge_file_pdf_celery: + queue: knowledge_pdf_celery + bisheng.worker.knowledge.*: # 其他知识库后台任务 + queue: celery bisheng.worker.workflow.*: # 工作流相关任务 queue: workflow_celery diff --git a/docs/dashboard-dataset-metric-calculation.md b/docs/dashboard-dataset-metric-calculation.md new file mode 100644 index 000000000..9c6275fec --- /dev/null +++ b/docs/dashboard-dataset-metric-calculation.md @@ -0,0 +1,130 @@ +# 看板数据集指标统计口径 + +> 更新时间:2026-08-05 +> 适用范围:`DASHBOARD_DATASET` 当前注册的全部数据集 +> 统计结果:14 个数据集、47 个指标 + +## 1. 文档目的 + +本文档说明看板中每个数据集及指标的端到端统计口径,包括: + +- 原始业务表或埋点事件来源; +- 中间表的记录生成粒度; +- 指标固定过滤条件; +- Elasticsearch 聚合方式; +- 最终计算公式; +- 时间范围和时间粒度; +- 指标去重口径。 + +指标定义来自 [init_dataset.py](../src/backend/bisheng/telemetry_search/domain/init_dataset.py#L15),查询执行逻辑位于 [component.py](../src/backend/bisheng/telemetry_search/domain/services/component.py#L155)。 + +## 2. 统一计算规则 + +| 规则 | 说明 | +|---|---| +| `cardinality(field)` | 按字段近似去重计数。 | +| `value_count(field)` | 统计字段非空的文档数量,不额外去重。 | +| 虚拟指标 | 聚合方式由后端数据集配置固定,看板组件中的聚合选项不改变其计算方式。 | +| 实体指标 | 默认使用 `sum`;看板编辑器允许改为平均、计数、最大、最小或去重计数。 | +| 比率指标 | 返回原始比值,例如 `0.25`;是否显示为 `25%` 由组件数字格式决定。分母为 `0` 时返回 `0`。 | +| 指标过滤 | 指标固定过滤条件与组件过滤、联动过滤、服务端权限过滤、时间过滤按 AND 合并。 | +| 累计总量,有时间维度 | 先按所选时间粒度计算每个时间桶的新增量,再执行 `cumulative_sum`。 | +| 累计总量,无时间维度 | 移除查询开始时间,保留结束时间,直接统计截至结束时间的总量。 | + +累计指标分支实现见 [component.py](../src/backend/bisheng/telemetry_search/domain/services/component.py#L260),ES 查询构造见 [search_engine_service.py](../src/backend/bisheng/telemetry_search/domain/services/search_engine_service.py#L338)。 + +## 3. 数据集数据来源与记录粒度 + +| 数据集 | 原始数据来源 | 中间表一条记录代表什么 | 时间字段口径 | 生成逻辑 | +|---|---|---|---|---| +| 用户行为指标表 `mid_user_increment` | MySQL `user` 表 | 一个用户;ES ID 为 `user_{user_id}` | 用户 `create_time` | [mid_table.py](../src/backend/bisheng/worker/telemetry/mid_table.py#L90) | +| 活跃用户表 `mid_active_user` | `base_telemetry_events` 中登录、会话、应用、知识库和文件操作事件 | 中国时区下“自然日+用户”一条记录,取当天该用户最新事件中的用户信息 | 最新一次活跃事件时间 | [derived_events.py](../src/backend/bisheng/telemetry/domain/mid_table/derived_events.py#L187) | +| 应用数量表 `mid_app_increment` | MySQL 应用/工作流数据 | 一个应用;ES ID 为 `app_{app_id}` | 应用 `create_time` | [mid_table.py](../src/backend/bisheng/worker/telemetry/mid_table.py#L866) | +| 会话数量表 `mid_sessions_increment` | `new_message_session` 埋点 | 一次新会话事件 | 埋点发生时间 | [derived_events.py](../src/backend/bisheng/telemetry/domain/mid_table/derived_events.py#L319) | +| 会话运行时长表 `mid_session_run_dtl` | `application_alive`、`application_process` 埋点 | 有运行时长的事件保存一条明细;时长为 0 的在线事件按覆盖的每一分钟展开一条记录 | 会话开始时间;并发使用 `minute_ts` | [derived_events.py](../src/backend/bisheng/telemetry/domain/mid_table/derived_events.py#L481) | +| 工具调用时长表 `mid_tool_call_dtl` | `tool_invoke` 埋点 | 一次工具调用事件 | 埋点发生时间 | [derived_events.py](../src/backend/bisheng/telemetry/domain/mid_table/derived_events.py#L359) | +| 知识库存量表 `mid_knowledge_increment` | MySQL `knowledge` 表 | 一个知识库;ES ID 为 `knowledge_{knowledge_id}` | 知识库 `create_time` | [mid_table.py](../src/backend/bisheng/worker/telemetry/mid_table.py#L931) | +| 知识库文件存量表 `mid_knowledge_file_increment` | MySQL `knowledgefile`、`qaknowledge` | 一个文档文件或一条 QA 数据;ES ID 分别为 `file-{id}`、`qa-{id}` | 文件或 QA 的 `create_time` | [knowledge_file_increment.py](../src/backend/bisheng/telemetry/domain/mid_table/knowledge_file_increment.py#L56) | +| 知识空间内容统计 `mid_knowledge_space_content_stat` | 当前成功、未删除、主版本的知识空间文件;文件预览事件 | 文件是一条当前快照;预览是“文件+中国自然日”一条累计记录 | 文件创建时间;预览日零点 | [mid_table.py](../src/backend/bisheng/worker/telemetry/mid_table.py#L372)、[knowledge_space_content.py](../src/backend/bisheng/telemetry/domain/mid_table/knowledge_space_content.py#L680) | +| 文件解析事件表 `mid_doc_parse_dtl` | `file_parse` 埋点 | 一次文件解析事件 | 埋点发生时间 | [derived_events.py](../src/backend/bisheng/telemetry/domain/mid_table/derived_events.py#L287) | +| 模型调用事件表 `mid_model_call_dtl` | `model_invoke` 埋点 | 一次模型调用按开始至结束覆盖的每一分钟展开;没有完整起止时间时只生成一条 | 埋点时间;并发使用 `minute_ts` | [derived_events.py](../src/backend/bisheng/telemetry/domain/mid_table/derived_events.py#L399) | +| 用户反馈指标表 `mid_user_interact_dtl` | `message_feedback` 埋点 | 一次点赞、点踩或复制操作 | 埋点发生时间 | [mid_table.py](../src/backend/bisheng/worker/telemetry/mid_table.py#L993) | +| 实时问答统计 `mid_realtime_qa_question_fact` | 专家问答 `Question`;门户 `portal_qa` 埋点 | “租户+问答类型+问题”一条记录 | 问题创建时间或问答成功时间 | [realtime_qa_question.py](../src/backend/bisheng/telemetry/domain/mid_table/realtime_qa_question.py#L91)、[realtime_dashboard.py](../src/backend/bisheng/worker/telemetry/realtime_dashboard.py#L119) | +| 全员每日参与度 `mid_user_daily_participation` | 当前有效员工名册+`user_login` 登录埋点 | “租户+中国自然日+用户”一条可更新记录 | 中国自然日零点;登录时间另存 | [daily_participation.py](../src/backend/bisheng/telemetry/domain/mid_table/daily_participation.py#L188)、[mid_table.py](../src/backend/bisheng/worker/telemetry/mid_table.py#L136) | + +## 4. 每个数据集下的指标计算 + +| 数据集 | 指标名称 | 指标字段 | 固定过滤条件 | ES 聚合方式 | 最终计算公式 | 时间口径 | 去重口径 | +|---|---|---|---|---|---|---|---| +| 用户行为指标表 | 总用户数 | `total_user_count` | 无 | 有时间维度:每桶 `cardinality(user_id)` 后累计;无时间维度:`cardinality(user_id)` | 截至当前时间桶的累计用户数 | 用户创建时间;支持年/月/周/日 | 按 `user_id` 去重 | +| 用户行为指标表 | 新增用户数 | `new_user_count` | 无 | `cardinality(user_id)` | 查询周期或时间桶内新增用户数 | 用户创建时间;支持年/月/周/日 | 按 `user_id` 去重 | +| 活跃用户表 | 活跃用户数 | `active_user_count` | 仅同步登录、会话、应用、知识库、知识文件相关活跃事件 | `cardinality(user_id)` | 查询周期或时间桶内至少产生一次活跃事件的用户数 | 活跃事件时间;支持年/月/周/日 | 按 `user_id` 去重;中间表已按日+用户合并 | +| 应用数量表 | 总应用数 | `total_app_count` | 无 | 有时间维度:每桶 `cardinality(app_id)` 后累计;无时间维度:`cardinality(app_id)` | 截至当前时间桶的累计应用数 | 应用创建时间;支持年/月/周/日 | 按 `app_id` 去重 | +| 应用数量表 | 新增应用数 | `new_app_count` | 无 | `cardinality(app_id)` | 查询周期或时间桶内新增应用数 | 应用创建时间;支持年/月/周/日 | 按 `app_id` 去重 | +| 会话数量表 | 会话次数 | `session_count` | 仅 `new_message_session` 事件 | `cardinality(session_id)` | 查询范围内不同会话的数量 | 会话创建事件时间;支持年/月/周/日 | 按 `session_id` 去重 | +| 会话数量表 | 使用人数 | `platform_user_count` | `source = platform` | `cardinality(user_id)` | 平台页面产生会话的不同用户数 | 会话创建事件时间;支持年/月/周/日 | 按 `user_id` 去重 | +| 会话数量表 | API 调用次数 | `api_call_count` | `source = api` | `cardinality(session_id)` | API 来源产生的不同会话数量 | 会话创建事件时间;支持年/月/周/日 | 按 `session_id` 去重 | +| 会话运行时长表 | 会话运行时长 | `duration_seconds` | `duration_seconds != 0` | 实体字段;默认 `sum(duration_seconds)`,可在组件中修改 | 默认等于符合条件记录的运行秒数之和 | 会话开始时间;支持年/月/周/日/小时 | 不去重,按运行事件明细求和 | +| 会话运行时长表 | 最大同时在线会话数 | `max_concurrent_sessions` | `duration_seconds = 0` | 按 `minute_ts` 创建 1 分钟桶;桶内 `cardinality(event_id)`;再执行 `max_bucket` | `max(每分钟不同在线事件数)` | 开始到结束覆盖的每一分钟,首尾分钟均包含 | 每分钟按 `event_id` 去重 | +| 工具调用时长表 | 工具调用次数 | `tool_call_count` | 无 | `value_count(event_id)` | 工具调用事件文档数 | 调用事件时间;支持年/月/周/日 | 不额外去重 | +| 工具调用时长表 | 工具调用成功次数 | `tool_call_success_count` | `status = success` | `value_count(event_id)` | 成功工具调用事件文档数 | 调用事件时间;支持年/月/周/日 | 不额外去重 | +| 工具调用时长表 | 工具调用成功率 | `tool_call_success_rate` | 分子:`status = success`;分母:全部事件 | 分子、分母分别执行 `value_count(event_id)` | 成功调用次数 ÷ 全部调用次数 | 调用事件时间;支持年/月/周/日 | 不额外去重 | +| 知识库存量表 | 总文档知识库数 | `total_document_knowledge_base_count` | `knowledge_type = 0` | 有时间维度:每桶 `cardinality(knowledge_id)` 后累计;无时间维度:直接去重统计 | 截至当前时间桶的文档知识库总量 | 知识库创建时间;支持年/月/周/日 | 按 `knowledge_id` 去重 | +| 知识库存量表 | 总 QA 知识库数 | `total_qa_knowledge_base_count` | `knowledge_type = 1` | 有时间维度:每桶去重后累计;无时间维度:直接去重统计 | 截至当前时间桶的 QA 知识库总量 | 知识库创建时间;支持年/月/周/日 | 按 `knowledge_id` 去重 | +| 知识库存量表 | 新增文档知识库数 | `new_document_knowledge_base_count` | `knowledge_type = 0` | `cardinality(knowledge_id)` | 查询周期或时间桶内新增文档知识库数 | 知识库创建时间;支持年/月/周/日 | 按 `knowledge_id` 去重 | +| 知识库存量表 | 新增 QA 知识库数 | `new_qa_knowledge_base_count` | `knowledge_type = 1` | `cardinality(knowledge_id)` | 查询周期或时间桶内新增 QA 知识库数 | 知识库创建时间;支持年/月/周/日 | 按 `knowledge_id` 去重 | +| 知识库文件存量表 | 总文件数 | `total_file_count` | `knowledge_base_type = 文档知识库` | 有时间维度:每桶 `value_count(file_id)` 后累计;无时间维度:`cardinality(file_id)` | 截至当前时间桶的文档文件总量 | 文件上传时间;支持年/月/周/日 | 无时间维度时按 `file_id` 去重 | +| 知识库文件存量表 | 总 QA 对数 | `total_qa_count` | `knowledge_base_type = QA知识库` | 有时间维度:每桶 `value_count(file_id)` 后累计;无时间维度:`cardinality(file_id)` | 截至当前时间桶的 QA 数据总量 | QA 创建时间;支持年/月/周/日 | 无时间维度时按 `file_id` 去重 | +| 知识库文件存量表 | 文件大小 | `file_size` | `knowledge_base_type = 文档知识库` | 实体字段;默认 `sum(file_size)`,可配置其他聚合 | 默认等于文档文件大小之和 | 文件上传时间;支持年/月/周/日 | 不去重,一条文件记录参与一次 | +| 知识空间内容统计 | 总文件数 | `total_file_count` | `record_type=file`、`file_type=1`、`space_level` 为 `public`、`department`、`team` 或 `team_ks` | 有时间维度:每桶 `value_count(file_id)` 后累计;无时间维度:直接 `value_count(file_id)` | 当前有效文件快照按创建时间形成的累计数量 | 文件创建时间;支持年/月/周/日 | 当前每个文件只有一条快照 | +| 知识空间内容统计 | 新增文件数 | `new_file_count` | 与“总文件数”相同 | `value_count(file_id)` | 查询周期或时间桶内创建、且当前仍有效的文件数量 | 文件创建时间;支持年/月/周/日 | 当前每个文件一条快照 | +| 知识空间内容统计 | 内容贡献人数 | `contributor_count` | `record_type=file`、空间级别在允许集合内 | `cardinality(uploader_user_id)` | 上传过当前有效内容的不同用户数 | 文件创建时间;支持年/月/周/日 | 按上传人 ID 去重 | +| 知识空间内容统计 | 预览次数 | `preview_count` | `record_type = preview_daily` | `sum(preview_count)` | 各“文件+自然日”记录中的预览计数之和 | 中国自然日零点;支持年/月/周/日 | 不按用户去重,每次成功预览均加 1 | +| 文件解析事件表 | 文档上传次数 | `doc_parse_count` | 无 | `value_count(event_id)` | 全部文件解析事件数量 | 解析事件时间;支持年/月/周/日 | 不额外去重 | +| 文件解析事件表 | 文档入库成功次数 | `doc_parse_success_count` | `status = success` | `value_count(event_id)` | 解析成功事件数量 | 解析事件时间;支持年/月/周/日 | 不额外去重 | +| 文件解析事件表 | 文档入库成功率 | `doc_parse_success_rate` | 分子:`status=success`;分母:全部解析事件 | 两次 `value_count(event_id)` | 解析成功次数 ÷ 全部解析次数 | 解析事件时间;支持年/月/周/日 | 不额外去重 | +| 文件解析事件表 | ETL 处理次数 | `etl_parse_count` | `parse_type = etl4lm` | `value_count(event_id)` | ETL 类型解析事件数量 | 解析事件时间;支持年/月/周/日 | 不额外去重 | +| 文件解析事件表 | ETL 处理成功次数 | `etl_parse_success_count` | `parse_type=etl4lm AND status=success` | `value_count(event_id)` | 成功的 ETL 解析事件数量 | 解析事件时间;支持年/月/周/日 | 不额外去重 | +| 文件解析事件表 | ETL 处理成功率 | `etl_parse_success_rate` | 分子:`parse_type=etl4lm AND status=success`;分母:`parse_type=etl4lm` | 两次 `value_count(event_id)` | ETL 成功次数 ÷ ETL 总次数 | 解析事件时间;支持年/月/周/日 | 不额外去重 | +| 模型调用事件表 | Token 消耗量 | `total_token` | 无 | 实体字段;默认 `sum(total_token)`,可配置其他聚合 | 默认等于命中文档中 `total_token` 之和 | 模型调用事件时间;支持年/月/周/日/小时 | 不去重;模型调用跨分钟时中间表会复制该字段 | +| 模型调用事件表 | 模型调用次数 | `model_call_count` | 无 | `cardinality(event_id)` | 不同模型调用事件数量 | 模型调用事件时间;支持年/月/周/日/小时 | 按 `event_id` 去重,消除分钟展开产生的重复 | +| 模型调用事件表 | 模型调用成功率 | `model_call_success_rate` | 分子:`status=success`;分母:全部模型调用 | 分子、分母分别 `cardinality(event_id)` | 成功模型调用数 ÷ 全部模型调用数 | 模型调用事件时间;支持年/月/周/日/小时 | 按 `event_id` 去重 | +| 模型调用事件表 | 最大 LLM 并发数 | `max_concurrent_llm_sessions` | 无 | 按 `minute_ts` 创建 1 分钟桶;桶内 `cardinality(event_id)`;再执行 `max_bucket` | `max(每分钟不同模型调用事件数)` | 调用开始至结束覆盖的每一分钟,首尾分钟均包含 | 每分钟按 `event_id` 去重 | +| 模型调用事件表 | 平均首 Token 响应延迟 | `avg_first_token_cost_time` | 无 | `avg(first_token_cost_time)` | 命中文档首 Token 延迟的算术平均值,单位毫秒 | 模型调用事件时间;支持年/月/周/日/小时 | 不按事件去重;按中间表文档计算平均 | +| 用户反馈指标表 | 点赞次数 | `like_count` | `interact_type = like` | `value_count(event_id)` | 点赞反馈事件数量 | 反馈事件时间;支持年/月/周/日 | 不额外去重 | +| 用户反馈指标表 | 点踩次数 | `dislike_count` | `interact_type = dislike` | `value_count(event_id)` | 点踩反馈事件数量 | 反馈事件时间;支持年/月/周/日 | 不额外去重 | +| 用户反馈指标表 | 复制次数 | `copy_count` | `interact_type = copy` | `value_count(event_id)` | 复制反馈事件数量 | 反馈事件时间;支持年/月/周/日 | 不额外去重 | +| 实时问答统计 | 问答总数 | `total_qa_count` | 无 | `value_count(question_id)` | 专家、智能、文档内 AI 问答记录总数 | 问题创建或成功时间;支持年/月/周/日/小时 | 中间表按租户+类型+问题 ID 保留一条 | +| 实时问答统计 | 专家问答数 | `expert_qa_count` | `qa_type = expert` | `value_count(question_id)` | 专家问答记录数 | 专家问题创建时间 | 每个专家问题一条记录 | +| 实时问答统计 | 智能问答数 | `smart_qa_count` | `qa_type = smart` | `value_count(question_id)` | 门户智能问答成功记录数 | 问答成功事件时间 | 每个问题一条记录 | +| 实时问答统计 | 文档内 AI 对话数 | `document_qa_count` | `qa_type = document` | `value_count(question_id)` | 文档场景问答成功记录数 | 问答成功事件时间 | 每个问题一条记录 | +| 实时问答统计 | 提问人数 | `qa_user_count` | 无 | `cardinality(user_id)` | 查询范围内至少产生一次问答的用户数 | 问题创建或成功时间;支持年/月/周/日/小时 | 按 `user_id` 去重 | +| 全员每日参与度 | 全员参与占比 | `participation_rate` | 分子:`logged_in=true`;分母:当天范围内全部名册记录 | 分子、分母分别 `value_count(user_id)` | 当天实际登录人数 ÷ 当天名册人数 | 中国自然日;默认查询当天 | 每天每租户每用户只有一条记录 | +| 全员每日参与度 | 实际登录人数 | `logged_in_employee_count` | `logged_in = true` | `value_count(user_id)` | 当天至少成功登录一次的员工数 | 中国自然日;默认查询当天 | 每天每个用户一条记录 | +| 全员每日参与度 | 全员总数 | `active_employee_count` | `active_employee = 1` | `value_count(user_id)` | 当天有效员工名册数量 | 中国自然日;默认查询当天 | 每天每个有效员工一条记录 | +| 全员每日参与度 | 实际登录次数 | `login_count` | 无 | 实体字段;默认 `sum(login_count)`,可配置其他聚合 | 成功登录事件的累计次数 | 中国自然日;默认查询当天 | 不按用户去重;同一用户每次成功登录均加 1 | + +## 5. 查询时间规则 + +| 场景 | 实际行为 | +|---|---| +| 普通数据集没有配置时间过滤 | 查询索引内全部数据。 | +| 实时问答、每日参与度没有配置时间过滤 | 自动使用当天 `00:00:00` 至 `23:59:59`。 | +| 同时存在组件时间过滤和联动时间过滤 | 取所有时间范围的交集。 | +| 使用时间维度 | 按看板选择的年、月、周、日或小时创建 `date_histogram`。 | +| 累计总量+时间维度 | 每个时间桶先统计新增量,再累计。 | +| 累计总量但没有时间维度 | 移除查询开始时间,只保留结束时间,统计截至结束时间的总量。 | + +时间过滤实现见 [component.py](../src/backend/bisheng/telemetry_search/domain/services/component.py#L464)。 + +## 6. 完整性校验 + +通过静态加载当前 `DASHBOARD_DATASET` 配置进行校验: + +```text +dataset_count=14 +metric_count=47 +unique_metric_keys=47 +``` + +本文档描述的是当前代码定义的统计计算逻辑,未连接实际数据库或 Elasticsearch 核验运行时数据。 diff --git a/features/v2.5.0-sg/release-contract.md b/features/v2.5.0-sg/release-contract.md index 68e06d975..ecc8cbbde 100644 --- a/features/v2.5.0-sg/release-contract.md +++ b/features/v2.5.0-sg/release-contract.md @@ -77,6 +77,7 @@ | INV-SG-23 | 新的知识文件链接/邀请码分享必须停止创建;上线前已创建的链接在原失效、密码、邀请码、下载和撤销规则下继续访问,并动态解析当前业务文档入口 | ShareLink, KnowledgeDocument, KnowledgeFile | F059 | | INV-SG-24 | canonical 内容变化与 `KnowledgeDocument.content_generation`、入口变化与对应 `KnowledgeFile` 的期望内容/入口代次必须在同一关系事务中提交;Worker 按入口期望/已应用双代次和 lease/CAS 幂等处理并恢复租户上下文,Celery 消息仅加速,定时任务必须补偿未完成或失败状态且不得伪报同步完成 | KnowledgeDocument, KnowledgeFile | F059 | | INV-SG-25 | F059 不扫描、推断、合并或回填上线前旧发布副本;新不变量只约束上线后由 F059 创建的逻辑入口关系 | KnowledgeDocument, KnowledgeFile | F059 | +| INV-SG-26 | `knowledge_celery` 只允许标题提取、首次文件解析和文件解析重试三个白名单任务;PDF Artifact 继续使用 `knowledge_pdf_celery`,工作流/审批继续使用 `workflow_celery`,其余后台任务进入默认 `celery` | Celery task routing | F060 | INV-SG-1 的“不继承”只约束推荐业务域特征,不修改基线 INV-12 中部门管理员的权限继承语义。 @@ -102,6 +103,7 @@ INV-SG-1 的“不继承”只约束推荐业务域特征,不修改基线 INV- | F059-knowledge-publish-share-unification | 既有 knowledge / version / tag / search / QA / stats 模块 | 依赖业务文档版本链、知识空间目录、MinIO、ES/Milvus、水印下载和统计能力 | | F059-knowledge-publish-share-unification | F004-rebac-core, F008-resource-rebac-adaptation | 依赖统一入口权限检查、`share_file` 和 OpenFGA 失败补偿 | | F059-knowledge-publish-share-unification | 既有 approval / share_link 模块 | 复用发布审批,新增固定两节点分享审批,并兼容已有知识文件分享链接 | +| F060-knowledge-parse-queue-isolation | 既有 Celery Worker 与 knowledge 文件解析任务 | 只调整任务路由所有权,不修改任务协议、解析实现或存量消息 | ```text F001 ──┐ @@ -140,6 +142,7 @@ approval/share_link ────────────────┘ | 2026-07-16 | 对齐远端业务域部门绑定实现,删除独立绑定领域对象和表,改用 `domains[].department_ids` 唯一配置源 | F056 | | 2026-07-19 | 登记 F057 对既有知识、权限和推荐投影的复用边界,并增加公共数据保护、判重与安全执行不变量 | F057 | | 2026-07-27 | 登记 F059 对 `KnowledgeFile` 逻辑入口及 `KnowledgeDocument` 状态驱动投影补偿的扩展边界;不新增领域表,并增加单实体发布、同级分享、权限、检索去重、容量和旧链接兼容不变量 | F059 | +| 2026-08-05 | 登记 F060 的 Celery 队列所有权:`knowledge_celery` 仅承载三个文件解析白名单任务 | F060 | --- diff --git a/src/backend/bisheng/api/services/knowledge_imp.py b/src/backend/bisheng/api/services/knowledge_imp.py index a7ce201e4..67975dab7 100644 --- a/src/backend/bisheng/api/services/knowledge_imp.py +++ b/src/backend/bisheng/api/services/knowledge_imp.py @@ -379,7 +379,7 @@ def addEmbedding( refresh_file_similarity_candidates_celery.apply_async( args=(db_file.id,), - queue="knowledge_celery", + queue="celery", ) except Exception: logger.exception("enqueue similarity candidate refresh failed file_id={}", db_file.id) diff --git a/src/backend/bisheng/config_3002.yaml b/src/backend/bisheng/config_3002.yaml index 6fb7d2068..616bab1cd 100644 --- a/src/backend/bisheng/config_3002.yaml +++ b/src/backend/bisheng/config_3002.yaml @@ -40,8 +40,16 @@ celery_task: knowledge_file_time_limit: null # 每秒处理10个文件 null: 表示无限制 # 对celery熟悉的用户可以自定义配置任务的路由,启动不同类型的worker处理不同类型的异步任务,注意工作流的执行只能在一个进程内 task_routers: - bisheng.worker.knowledge.*: # 知识库文件处理相关任务 + bisheng.worker.knowledge.file_title_worker.extract_knowledge_file_title_celery: queue: knowledge_celery + bisheng.worker.knowledge.file_worker.parse_knowledge_file_celery: + queue: knowledge_celery + bisheng.worker.knowledge.file_worker.retry_knowledge_file_celery: + queue: knowledge_celery + bisheng.worker.knowledge.pdf_artifact_worker.generate_knowledge_file_pdf_celery: + queue: knowledge_pdf_celery + bisheng.worker.knowledge.*: # 其他知识库后台任务 + queue: celery bisheng.worker.workflow.*: # 工作流相关任务 queue: workflow_celery diff --git a/src/backend/bisheng/config_3003.yaml b/src/backend/bisheng/config_3003.yaml index eb7fc6559..aec07fb3c 100644 --- a/src/backend/bisheng/config_3003.yaml +++ b/src/backend/bisheng/config_3003.yaml @@ -40,8 +40,16 @@ celery_task: knowledge_file_time_limit: null # 每秒处理10个文件 null: 表示无限制 # 对celery熟悉的用户可以自定义配置任务的路由,启动不同类型的worker处理不同类型的异步任务,注意工作流的执行只能在一个进程内 task_routers: - bisheng.worker.knowledge.*: # 知识库文件处理相关任务 + bisheng.worker.knowledge.file_title_worker.extract_knowledge_file_title_celery: queue: knowledge_celery + bisheng.worker.knowledge.file_worker.parse_knowledge_file_celery: + queue: knowledge_celery + bisheng.worker.knowledge.file_worker.retry_knowledge_file_celery: + queue: knowledge_celery + bisheng.worker.knowledge.pdf_artifact_worker.generate_knowledge_file_pdf_celery: + queue: knowledge_pdf_celery + bisheng.worker.knowledge.*: # 其他知识库后台任务 + queue: celery bisheng.worker.workflow.*: # 工作流相关任务 queue: workflow_celery diff --git a/src/backend/bisheng/core/config/celery_queues.py b/src/backend/bisheng/core/config/celery_queues.py new file mode 100644 index 000000000..9cdcad2d1 --- /dev/null +++ b/src/backend/bisheng/core/config/celery_queues.py @@ -0,0 +1,61 @@ +"""Celery queue ownership and routing contracts.""" + +from __future__ import annotations + +from collections.abc import Mapping +from typing import Any + +DEFAULT_CELERY_QUEUE = "celery" +KNOWLEDGE_PARSE_QUEUE = "knowledge_celery" +KNOWLEDGE_PDF_QUEUE = "knowledge_pdf_celery" +WORKFLOW_CELERY_QUEUE = "workflow_celery" + +KNOWLEDGE_PARSE_TASKS = frozenset( + { + "bisheng.worker.knowledge.file_title_worker.extract_knowledge_file_title_celery", + "bisheng.worker.knowledge.file_worker.parse_knowledge_file_celery", + "bisheng.worker.knowledge.file_worker.retry_knowledge_file_celery", + } +) +PDF_ARTIFACT_TASK = "bisheng.worker.knowledge.pdf_artifact_worker.generate_knowledge_file_pdf_celery" + +_DEFAULT_QUEUE_PATTERNS = ( + "bisheng.worker.knowledge.*", + "bisheng.worker.org_sync.*", + "bisheng.worker.tenant_reconcile.*", + "bisheng.worker.admin_scope.*", + "bisheng.worker.message.*", + "bisheng.worker.portal_course.*", + "bisheng.worker.permission.*", +) +_WORKFLOW_QUEUE_PATTERNS = ( + "bisheng.worker.workflow.*", + "bisheng.worker.approval.*", +) + + +def _normalize_configured_route(route: Any) -> Any: + """Prevent legacy custom routes from assigning work to the parse queue.""" + if isinstance(route, str): + return DEFAULT_CELERY_QUEUE if route == KNOWLEDGE_PARSE_QUEUE else route + if isinstance(route, Mapping): + normalized = dict(route) + if normalized.get("queue") == KNOWLEDGE_PARSE_QUEUE: + normalized["queue"] = DEFAULT_CELERY_QUEUE + return normalized + return route + + +def build_celery_task_routes(configured_routes: Mapping[str, Any] | None) -> dict[str, Any]: + """Build final routes while enforcing exclusive parse-queue ownership.""" + routes: dict[str, Any] = {pattern: {"queue": WORKFLOW_CELERY_QUEUE} for pattern in _WORKFLOW_QUEUE_PATTERNS} + routes.update({pattern: {"queue": DEFAULT_CELERY_QUEUE} for pattern in _DEFAULT_QUEUE_PATTERNS}) + + for task_pattern, route in (configured_routes or {}).items(): + if task_pattern in routes or task_pattern in KNOWLEDGE_PARSE_TASKS: + continue + routes[task_pattern] = _normalize_configured_route(route) + + routes[PDF_ARTIFACT_TASK] = {"queue": KNOWLEDGE_PDF_QUEUE} + routes.update({task_name: {"queue": KNOWLEDGE_PARSE_QUEUE} for task_name in sorted(KNOWLEDGE_PARSE_TASKS)}) + return routes diff --git a/src/backend/bisheng/core/config/settings.py b/src/backend/bisheng/core/config/settings.py index 9020391ef..37c1b692d 100644 --- a/src/backend/bisheng/core/config/settings.py +++ b/src/backend/bisheng/core/config/settings.py @@ -9,6 +9,7 @@ from cryptography.fernet import Fernet from loguru import logger from pydantic import BaseModel, ConfigDict, Field, field_validator, model_validator +from bisheng.core.config.celery_queues import build_celery_task_routes from bisheng.core.config.llm import LLMConf from bisheng.core.config.multi_tenant import MultiTenantConf from bisheng.core.config.openfga import OpenFGAConf @@ -162,24 +163,7 @@ class CeleryConf(BaseModel): @model_validator(mode="after") def validate(self): - if not self.task_routers: - self.task_routers = { - "bisheng.worker.knowledge.*": {"queue": "knowledge_celery"}, # Knowledge Base Related Tasks - "bisheng.worker.workflow.*": {"queue": "workflow_celery"}, # Workflow Execution Related Tasks - "bisheng.worker.org_sync.*": {"queue": "knowledge_celery"}, - # Org Sync Tasks (low frequency, reuse knowledge queue) - "bisheng.worker.tenant_reconcile.*": {"queue": "knowledge_celery"}, - # v2.5.1 F012 — 6h catch-up, reuse knowledge_celery - "bisheng.worker.admin_scope.*": {"queue": "knowledge_celery"}, # v2.5.1 F019 — 10min sweep, low-volume - "bisheng.worker.message.*": {"queue": "knowledge_celery"}, # WeChat message push tasks - "bisheng.worker.portal_course.*": {"queue": "knowledge_celery"}, - "bisheng.worker.permission.*": {"queue": "knowledge_celery"}, - } - else: - self.task_routers.setdefault( - "bisheng.worker.permission.*", - {"queue": "knowledge_celery"}, - ) + self.task_routers = build_celery_task_routes(self.task_routers) if "telemetry_mid_user_increment" not in self.beat_schedule: self.beat_schedule["telemetry_mid_user_increment"] = { "task": "bisheng.worker.telemetry.mid_table.sync_mid_user_increment", diff --git a/src/backend/bisheng/knowledge/domain/services/favorite_notify.py b/src/backend/bisheng/knowledge/domain/services/favorite_notify.py index 45538ebc0..13c0957cb 100644 --- a/src/backend/bisheng/knowledge/domain/services/favorite_notify.py +++ b/src/backend/bisheng/knowledge/domain/services/favorite_notify.py @@ -239,7 +239,7 @@ def enqueue_favorite_change_events(events: Iterable[FavoriteChangeEvent]) -> Non try: send_favorite_change_notifications.apply_async( args=[[event.model_dump(mode="json") for event in batch]], - queue="knowledge_celery", + queue="celery", ) except Exception: # 通知是明确的尽力流程;入队失败不能回滚已完成的文档操作。 diff --git a/src/backend/bisheng/knowledge/domain/services/knowledge_migration_service.py b/src/backend/bisheng/knowledge/domain/services/knowledge_migration_service.py index c72f616bf..a845582e3 100644 --- a/src/backend/bisheng/knowledge/domain/services/knowledge_migration_service.py +++ b/src/backend/bisheng/knowledge/domain/services/knowledge_migration_service.py @@ -63,7 +63,7 @@ class CeleryKnowledgeMigrationTaskDispatcher: task = preflight_knowledge_migration.apply_async( args=[batch_id], - queue="knowledge_celery", + queue="celery", ) return str(task.id) if task.id else None @@ -72,7 +72,7 @@ class CeleryKnowledgeMigrationTaskDispatcher: task = execute_knowledge_migration.apply_async( args=[batch_id, round_no], - queue="knowledge_celery", + queue="celery", ) return str(task.id) if task.id else None diff --git a/src/backend/bisheng/knowledge/domain/services/knowledge_space_service.py b/src/backend/bisheng/knowledge/domain/services/knowledge_space_service.py index d7deb14e4..dc9bbffa8 100644 --- a/src/backend/bisheng/knowledge/domain/services/knowledge_space_service.py +++ b/src/backend/bisheng/knowledge/domain/services/knowledge_space_service.py @@ -5714,7 +5714,7 @@ class KnowledgeSpaceService(KnowledgeUtils): "searched_at": searched_at.isoformat() if isinstance(searched_at, datetime) else searched_at, }, headers={"tenant_id": int(payload["tenant_id"])}, - queue="knowledge_celery", + queue="celery", expires=600, ) diff --git a/src/backend/bisheng/knowledge/domain/services/portal_hot_search_admin_service.py b/src/backend/bisheng/knowledge/domain/services/portal_hot_search_admin_service.py index 14406e54a..65ff07677 100644 --- a/src/backend/bisheng/knowledge/domain/services/portal_hot_search_admin_service.py +++ b/src/backend/bisheng/knowledge/domain/services/portal_hot_search_admin_service.py @@ -17,7 +17,7 @@ from bisheng.knowledge.domain.schemas.portal_hot_search_schema import ( ) from bisheng.utils.http_middleware import _check_is_global_super -KNOWLEDGE_QUEUE = "knowledge_celery" +DEFAULT_QUEUE = "celery" class PortalHotSearchAdminService: @@ -64,7 +64,7 @@ class PortalHotSearchAdminService: async_result = trigger_portal_hot_search_rebuild_celery.apply_async( headers={"tenant_id": tenant_id}, - queue=KNOWLEDGE_QUEUE, + queue=DEFAULT_QUEUE, ) return str(async_result.id) @@ -72,5 +72,5 @@ class PortalHotSearchAdminService: def _dispatch_fanout_rebuild() -> str: from bisheng.worker.knowledge.portal_hot_search import fanout_portal_hot_search_rebuild - async_result = fanout_portal_hot_search_rebuild.apply_async(queue=KNOWLEDGE_QUEUE) + async_result = fanout_portal_hot_search_rebuild.apply_async(queue=DEFAULT_QUEUE) return str(async_result.id) diff --git a/src/backend/bisheng/org_sync/api/endpoints/sync_exec.py b/src/backend/bisheng/org_sync/api/endpoints/sync_exec.py index 030e73e08..712580d33 100644 --- a/src/backend/bisheng/org_sync/api/endpoints/sync_exec.py +++ b/src/backend/bisheng/org_sync/api/endpoints/sync_exec.py @@ -67,7 +67,7 @@ async def execute_sync( from bisheng.worker.org_sync.tasks import execute_org_sync execute_org_sync.apply_async( args=[config_id, 'manual', login_user.user_id], - queue='knowledge_celery', + queue='celery', ) return resp_200({ diff --git a/src/backend/bisheng/permission/domain/services/primary_department_change_coordinator.py b/src/backend/bisheng/permission/domain/services/primary_department_change_coordinator.py index c00b8a9e3..6c8ce327d 100644 --- a/src/backend/bisheng/permission/domain/services/primary_department_change_coordinator.py +++ b/src/backend/bisheng/permission/domain/services/primary_department_change_coordinator.py @@ -291,4 +291,4 @@ async def cancel_primary_department_change( async def _dispatch_cleanup_event(event_id: int) -> None: from bisheng.worker.permission.department_transfer_cleanup import process_event - process_event.apply_async(args=[event_id], queue="knowledge_celery") + process_event.apply_async(args=[event_id], queue="celery") diff --git a/src/backend/bisheng/telemetry_search/domain/init_dataset.py b/src/backend/bisheng/telemetry_search/domain/init_dataset.py index 895f5f1b4..fcdd5d886 100644 --- a/src/backend/bisheng/telemetry_search/domain/init_dataset.py +++ b/src/backend/bisheng/telemetry_search/domain/init_dataset.py @@ -491,12 +491,16 @@ DASHBOARD_DATASET = [ field="tool_type" ), DimensionConfig( - name="应用名称", + name="应用类型", field="app_type" ), DimensionConfig( name="应用ID", field="app_id" + ), + DimensionConfig( + name="应用名称", + field="app_name" ) ] ).model_dump() @@ -1049,7 +1053,7 @@ DASHBOARD_DATASET = [ field="status" ), DimensionConfig( - name="文件来源", + name="应用类型", field="app_type" ) ] @@ -1585,23 +1589,25 @@ def _clone_dashboard_dataset(dataset: DashboardDataset) -> DashboardDataset: return DashboardDataset(**data) -REALTIME_DASHBOARD_DATASET_CODES = ( +DASHBOARD_DATASET_REFRESH_CODES = ( "mid_knowledge_space_content_stat", "mid_realtime_qa_question_fact", "mid_user_daily_participation", + "mid_tool_call_dtl", + "mid_doc_parse_dtl", ) -async def _upgrade_realtime_dashboard_datasets( +async def _upgrade_dashboard_datasets( dashboard_dataset_repository: DashboardDatasetRepositoryImpl, ): - """Idempotently install or refresh the three real-time dashboard datasets.""" + """Idempotently install or refresh dashboard datasets with schema updates.""" seed_by_code = { dataset.dataset_code: dataset for dataset in DASHBOARD_DATASET - if dataset.dataset_code in REALTIME_DASHBOARD_DATASET_CODES + if dataset.dataset_code in DASHBOARD_DATASET_REFRESH_CODES } - for dataset_code in REALTIME_DASHBOARD_DATASET_CODES: + for dataset_code in DASHBOARD_DATASET_REFRESH_CODES: seed_dataset = seed_by_code[dataset_code] existing_dataset = await dashboard_dataset_repository.find_one( dataset_code=dataset_code @@ -1630,7 +1636,7 @@ async def init_dashboard_datasets(): else: # Upgrade path: add department dimensions to existing datasets await _upgrade_datasets_add_department_dimensions(dashboard_dataset_repository) - await _upgrade_realtime_dashboard_datasets(dashboard_dataset_repository) + await _upgrade_dashboard_datasets(dashboard_dataset_repository) preset_dashboard = await DashboardDao.get_dashboards(dashboard_type=[DashboardType.PRESET_OSS]) if not preset_dashboard: await DashboardDao.exec_sql_str(preset_oss_dashboard_sql) diff --git a/src/backend/bisheng/telemetry_search/domain/services/dashboard.py b/src/backend/bisheng/telemetry_search/domain/services/dashboard.py index 72da4398a..9ee29fe23 100644 --- a/src/backend/bisheng/telemetry_search/domain/services/dashboard.py +++ b/src/backend/bisheng/telemetry_search/domain/services/dashboard.py @@ -44,6 +44,112 @@ class DashboardService(BaseModel): "department": "部门库", "team": "团队库(含科室库)", } + APPLICATION_TYPE_LABELS: ClassVar[dict[str, str]] = { + "workflow": "工作流", + "assistant": "助手", + "linsight": "Linsight", + "daily_chat": "日常对话", + "knowledge_base": "知识库", + "knowledge_space": "知识空间", + "rag_traceability": "RAG溯源", + "evaluation": "模型评测", + "model_test": "模型连通性测试", + "asr": "语音识别", + "tts": "语音合成", + "unknown": "未知", + } + STATUS_LABELS: ClassVar[dict[str, str]] = { + "success": "成功", + "failed": "失败", + "parse_failed": "解析失败", + } + DASHBOARD_ENUM_LABELS: ClassVar[ + dict[str, dict[str, dict[str, str]]] + ] = { + "mid_app_increment": { + "app_type": APPLICATION_TYPE_LABELS, + }, + "mid_sessions_increment": { + "source": { + "platform": "平台端", + "api": "API调用", + }, + "app_id": APPLICATION_TYPE_LABELS, + }, + "mid_session_run_dtl": { + "app_id": APPLICATION_TYPE_LABELS, + }, + "mid_tool_call_dtl": { + "tool_type": { + "0": "API工具", + "1": "内置工具", + "2": "MCP工具", + }, + "app_type": APPLICATION_TYPE_LABELS, + "app_id": APPLICATION_TYPE_LABELS, + }, + "mid_doc_parse_dtl": { + "parse_type": { + "local": "本地解析", + "uns": "UNS解析", + "etl4lm": "ETL4LM解析", + "un_etl4lm": "非ETL4LM解析", + "mineru": "MinerU解析", + "paddle_ocr": "PaddleOCR解析", + }, + "status": STATUS_LABELS, + "app_type": APPLICATION_TYPE_LABELS, + }, + "mid_model_call_dtl": { + "model_type": { + "llm": "大语言模型", + "embedding": "嵌入模型", + "rerank": "重排模型", + "asr": "语音识别", + "tts": "语音合成", + }, + "app_id": APPLICATION_TYPE_LABELS, + }, + "mid_realtime_qa_question_fact": { + "department_source": { + "event_time": "提问时所属主部门", + "current_primary_backfill": "当前主部门(历史回填)", + }, + "scene": { + "expert_question": "专家问答", + "smart_qa": "智能问答", + "document_qa": "知识门户·文档问答", + "my_knowledge_document_qa": "我的知识·文档问答", + }, + "source_app": { + "bisheng_my_knowledge": "毕昇·我的知识", + "expert_qa": "专家问答", + "shougang_portal": "首钢知识门户", + }, + }, + "mid_user_daily_participation": { + "department_source": { + "event_time": "登录时所属主部门", + "current_roster": "当前在职名册", + "current_roster_backfill": "当前名册(历史回填)", + "current_primary_backfill": "当前主部门(历史登录回填)", + }, + }, + } + + @staticmethod + def _format_date_enum_label(value: Any) -> Any: + if not isinstance(value, (int, float)) or isinstance(value, bool): + return value + return datetime.fromtimestamp(value / 1000).strftime("%Y-%m-%d %H:%M:%S") + + @classmethod + def _get_enum_labels( + cls, + dataset_code: str, + field: str, + ) -> dict[str, str]: + return cls.DASHBOARD_ENUM_LABELS.get(dataset_code, {}).get(field, {}) @classmethod def _uses_realtime_dataset( @@ -582,15 +688,18 @@ class DashboardService(BaseModel): if not dataset: raise NotFoundError() schema_config = dataset.schema_config if isinstance(dataset.schema_config, dict) else {} - allowed_fields = { - dimension.get("field") + dimension_by_field = { + dimension.get("field"): dimension for dimension in schema_config.get("dimensions", []) if dimension.get("field") } + allowed_fields = set(dimension_by_field) if field not in allowed_fields or ( label_field and label_field not in allowed_fields ): raise UnAuthorizedError() + label_dimension = dimension_by_field[label_field or field] + label_field_type = label_dimension.get("field_type") or label_dimension.get("type") scope_filters = await self._get_realtime_scope_filters(dataset_code) skip = (page - 1) * size @@ -634,9 +743,29 @@ class DashboardService(BaseModel): current_aggs = core_aggs + enum_labels = self._get_enum_labels(dataset_code, field) if keyword: search_field = label_field or field - filter_query = {"match_phrase": {f"{search_field}.text": keyword}} + text_filter = {"match_phrase": {f"{search_field}.text": keyword}} + if enum_labels: + normalized_keyword = keyword.casefold() + matched_values = [ + value + for value, label in enum_labels.items() + if normalized_keyword in value.casefold() + or normalized_keyword in label.casefold() + ] + search_filters = [text_filter] + if matched_values: + search_filters.append({"terms": {field: matched_values}}) + filter_query = { + "bool": { + "should": search_filters, + "minimum_should_match": 1, + } + } + else: + filter_query = text_filter current_aggs = { "filter_wrapper": { @@ -704,6 +833,10 @@ class DashboardService(BaseModel): if label_buckets else value ) + if label_field_type == "date": + label = self._format_date_enum_label(label) + else: + label = enum_labels.get(str(value), label) options.append({"value": value, "label": label}) if ( dataset_code == "mid_knowledge_space_content_stat" diff --git a/src/backend/bisheng/worker/config.py b/src/backend/bisheng/worker/config.py index ebc6066f3..8d3d02eb8 100644 --- a/src/backend/bisheng/worker/config.py +++ b/src/backend/bisheng/worker/config.py @@ -1,4 +1,5 @@ from bisheng.common.services.config_service import settings +from bisheng.core.config.celery_queues import build_celery_task_routes from bisheng.core.config.celery_redis import build_celery_redis_config _celery_redis_config = build_celery_redis_config(settings.celery_redis_url) @@ -11,13 +12,7 @@ result_serializer = "json" accept_content = ["json"] timezone = "Asia/Shanghai" enable_utc = False -_DEFAULT_ROUTES = { - "bisheng.worker.knowledge.pdf_artifact_worker.generate_knowledge_file_pdf_celery": { - "queue": "knowledge_pdf_celery" - }, - "bisheng.worker.approval.*": {"queue": "workflow_celery"}, -} -task_routes = {**_DEFAULT_ROUTES, **settings.celery_task.task_routers} +task_routes = build_celery_task_routes(settings.celery_task.task_routers) # redisHealth check interval, unit sec redis_backend_health_check_interval = 5 diff --git a/src/backend/bisheng/worker/knowledge/document_projection.py b/src/backend/bisheng/worker/knowledge/document_projection.py index 2ff2dbbce..04e705bdd 100644 --- a/src/backend/bisheng/worker/knowledge/document_projection.py +++ b/src/backend/bisheng/worker/knowledge/document_projection.py @@ -56,7 +56,7 @@ from bisheng.worker.approval.tasks import ( from bisheng.worker.main import bisheng_celery logger = logging.getLogger(__name__) -KNOWLEDGE_QUEUE = "knowledge_celery" +DEFAULT_QUEUE = "celery" SCAN_PAGE_SIZE = 100 @@ -547,7 +547,7 @@ async def _reconcile_rollback_candidates( ), }, headers={"tenant_id": int(tenant_id)}, - queue=KNOWLEDGE_QUEUE, + queue=DEFAULT_QUEUE, ) dispatched += 1 return dispatched @@ -610,7 +610,7 @@ async def _scan_tenant_projection_async(tenant_id: int) -> int: "entry_id": int(entry_id), }, headers={"tenant_id": int(tenant_id)}, - queue=KNOWLEDGE_QUEUE, + queue=DEFAULT_QUEUE, ) await _reconcile_permission_candidates( tenant_id=tenant_id, @@ -656,7 +656,7 @@ async def _fanout_projection_scan_async() -> int: scan_tenant_document_projections.apply_async( kwargs={"tenant_id": tenant_id}, headers={"tenant_id": tenant_id}, - queue=KNOWLEDGE_QUEUE, + queue=DEFAULT_QUEUE, ) return len(set(tenant_ids)) @@ -680,7 +680,7 @@ def enqueue_document_projection_entries( scan_tenant_document_projections.apply_async( kwargs={"tenant_id": int(tenant_id)}, headers={"tenant_id": int(tenant_id)}, - queue=KNOWLEDGE_QUEUE, + queue=DEFAULT_QUEUE, ) return for entry_id in sorted({int(item) for item in entry_ids}): @@ -690,7 +690,7 @@ def enqueue_document_projection_entries( "entry_id": entry_id, }, headers={"tenant_id": int(tenant_id)}, - queue=KNOWLEDGE_QUEUE, + queue=DEFAULT_QUEUE, ) @@ -797,5 +797,5 @@ def enqueue_knowledge_space_retirement(*, tenant_id: int, space_id: int) -> None process_knowledge_space_retirement.apply_async( kwargs={"tenant_id": int(tenant_id), "space_id": int(space_id)}, headers={"tenant_id": int(tenant_id)}, - queue=KNOWLEDGE_QUEUE, + queue=DEFAULT_QUEUE, ) diff --git a/src/backend/bisheng/worker/knowledge/portal_hot_search.py b/src/backend/bisheng/worker/knowledge/portal_hot_search.py index e48f60bba..44f54a043 100644 --- a/src/backend/bisheng/worker/knowledge/portal_hot_search.py +++ b/src/backend/bisheng/worker/knowledge/portal_hot_search.py @@ -43,12 +43,12 @@ from bisheng.knowledge.domain.services.portal_hot_search_scoring_service import from bisheng.worker._asyncio_utils import run_async_task from bisheng.worker.main import bisheng_celery -KNOWLEDGE_QUEUE = "knowledge_celery" +DEFAULT_QUEUE = "celery" def _dispatch_task_for_tenants(task, tenant_ids: list[int]) -> None: for tenant_id in sorted({int(value) for value in tenant_ids if int(value) > 0}): - task.apply_async(headers={"tenant_id": tenant_id}, queue=KNOWLEDGE_QUEUE) + task.apply_async(headers={"tenant_id": tenant_id}, queue=DEFAULT_QUEUE) def _build_llm_invoke(tenant_id: int) -> Callable[[str], str] | None: diff --git a/src/backend/bisheng/worker/knowledge/portal_recommendation.py b/src/backend/bisheng/worker/knowledge/portal_recommendation.py index 05c6f88e6..83f2ccf7a 100644 --- a/src/backend/bisheng/worker/knowledge/portal_recommendation.py +++ b/src/backend/bisheng/worker/knowledge/portal_recommendation.py @@ -49,7 +49,7 @@ from bisheng.shougang_portal_config.domain.services.portal_config_service import from bisheng.worker._asyncio_utils import run_async_task from bisheng.worker.main import bisheng_celery -KNOWLEDGE_QUEUE = "knowledge_celery" +DEFAULT_QUEUE = "celery" PROJECTION_PAGE_SIZE = 500 @@ -70,7 +70,7 @@ def _dispatch_task_for_tenants(task, tenant_ids: list[int], *, kwargs: dict | No task.apply_async( kwargs=dict(kwargs or {}), headers={"tenant_id": tenant_id}, - queue=KNOWLEDGE_QUEUE, + queue=DEFAULT_QUEUE, ) @@ -91,7 +91,7 @@ def enqueue_portal_recommendation_projection_refresh( "deleted": bool(deleted), }, headers={"tenant_id": resolved_tenant_id}, - queue=KNOWLEDGE_QUEUE, + queue=DEFAULT_QUEUE, ) @@ -113,7 +113,7 @@ def enqueue_portal_recommendation_projection_refresh_batch( "deleted": bool(deleted), }, headers={"tenant_id": resolved_tenant_id}, - queue=KNOWLEDGE_QUEUE, + queue=DEFAULT_QUEUE, ) @@ -132,7 +132,7 @@ def enqueue_portal_recommendation_resource_refresh( "event_version": event_version, }, headers={"tenant_id": resolved_tenant_id}, - queue=KNOWLEDGE_QUEUE, + queue=DEFAULT_QUEUE, ) @@ -148,7 +148,7 @@ def enqueue_portal_recommendation_user_invalidation( invalidate_user_ids_celery.apply_async( kwargs={"user_ids": normalized_ids}, headers={"tenant_id": resolved_tenant_id}, - queue=KNOWLEDGE_QUEUE, + queue=DEFAULT_QUEUE, ) @@ -163,12 +163,12 @@ def enqueue_portal_recommendation_config_post_commit( invalidate_department_users_celery.apply_async( kwargs={"department_ids": sorted({int(value) for value in department_ids})}, headers=headers, - queue=KNOWLEDGE_QUEUE, + queue=DEFAULT_QUEUE, ) if rebuild_pools: prepare_pool_rebuild_celery.apply_async( headers=headers, - queue=KNOWLEDGE_QUEUE, + queue=DEFAULT_QUEUE, ) @@ -176,7 +176,7 @@ def enqueue_portal_recommendation_pool_rebuild(*, tenant_id: int | None = None) resolved_tenant_id = int(tenant_id or get_current_tenant_id() or DEFAULT_TENANT_ID) prepare_pool_rebuild_celery.apply_async( headers={"tenant_id": resolved_tenant_id}, - queue=KNOWLEDGE_QUEUE, + queue=DEFAULT_QUEUE, ) @@ -794,7 +794,7 @@ async def _prepare_pool_rebuild_async() -> bool: "fingerprint": fingerprint, }, headers={"tenant_id": tenant_id}, - queue=KNOWLEDGE_QUEUE, + queue=DEFAULT_QUEUE, ) return True diff --git a/src/backend/bisheng/worker/org_sync/reconcile_tasks.py b/src/backend/bisheng/worker/org_sync/reconcile_tasks.py index 9e2b1e3c7..3b2037343 100644 --- a/src/backend/bisheng/worker/org_sync/reconcile_tasks.py +++ b/src/backend/bisheng/worker/org_sync/reconcile_tasks.py @@ -52,7 +52,7 @@ async def _fan_out_all() -> None: continue try: reconcile_single_config.apply_async( - args=[c.id], queue='knowledge_celery', + args=[c.id], queue='celery', ) except Exception as e: logger.exception( diff --git a/src/backend/bisheng/worker/org_sync/tasks.py b/src/backend/bisheng/worker/org_sync/tasks.py index a98105620..3f04da668 100644 --- a/src/backend/bisheng/worker/org_sync/tasks.py +++ b/src/backend/bisheng/worker/org_sync/tasks.py @@ -18,7 +18,7 @@ logger = logging.getLogger(__name__) @bisheng_celery.task(acks_late=True, time_limit=1800, soft_time_limit=1500) def execute_org_sync(config_id: int, trigger_type: str, trigger_user: int = None): - """Execute org sync for a given config. Runs in knowledge_celery queue.""" + """Execute org sync for a given config. Runs in the default queue.""" run_async_task(lambda: _execute_org_sync_async(config_id, trigger_type, trigger_user)) @@ -70,7 +70,7 @@ async def _check_schedules_async(): ) execute_org_sync.apply_async( args=[config.id, 'scheduled', None], - queue='knowledge_celery', + queue='celery', ) except Exception as e: logger.warning( diff --git a/src/backend/bisheng/worker/permission/department_transfer_cleanup.py b/src/backend/bisheng/worker/permission/department_transfer_cleanup.py index 741030ef5..7471e6c41 100644 --- a/src/backend/bisheng/worker/permission/department_transfer_cleanup.py +++ b/src/backend/bisheng/worker/permission/department_transfer_cleanup.py @@ -103,7 +103,7 @@ async def _scan_due_events_async() -> int: await _mark_overdue_if_needed(repository, event, now=now) await session.commit() try: - process_event.apply_async(args=[event_id], queue="knowledge_celery") + process_event.apply_async(args=[event_id], queue="celery") dispatched += 1 except Exception as exc: logger.warning( diff --git a/src/backend/celerybeat-schedule.db b/src/backend/celerybeat-schedule.db index 57e2dc96a137aeedbfe0e559c592ef94ae242062..2599af5ad0be81dfe0f68dd5ab5e0fc3448ff02b 100644 GIT binary patch delta 1168 zcmeH@y-!nd6o!9oZ@INqZd*lLq>6%46ahv27ElmBP>>Hrln*Q3q$F2O6(N7L`Gyqwb(6c zL{97xdqu6-C-#dvQ7;aNgQ7v?MWZ+*4vQvnL^O*Q(JGFLW1>wQ7wzJNI4Mqv)8dRc zE6$1YqC<3wE^$Fz6y2gnToS$FvbZ9yiff`zTo*S)zZeiV#h@4x!(v3-61T;uxFg2I zU2#v0iwQ9)ro??w5D&z(c=#2!qa~KLBgV0_wjUJe*#7dU>)Yw}d3VNr>dw!;oSJo~ zpS^Hry~VNSTzYE8EzG%d<)W2;ycRO%UB9<7F53cPyGtrxLAi%EhX; z<)Lz)DOTr;Ij_*0$$CFKL)q%!RWxkOKNL48J(J$zK&H(Jr030H^JLJNRmUX#c$d5C z)L3TCiI_)&mByG#Q?fpN&jeR6ugZCExzWx0JzefM$4@1#700hlTE8~A*}!@OOQAyC U50?_w(wB--;(y@>?w|L60i3#}x&QzG delta 1039 zcmeIuIZu-T0LJmR7sXb2TNR|9c;W%Aii#(ocwg9x_lf!rYU1EAn$+Pxt4U2)2NO3p z^Ufy!572>ui5Le5Hy09}tcxGOnc;c1=NW$W_cN(_J8$|k)gVZHcsnqjs5OhN#N;b6 zfk{kZ8Z+2|S?t6Gn8Q3S#6{SJAr`QRi*X4q#cu4uW!Q@)T#hSnC9cBNxCZ-hEw01$ zxB)lfChW(}xCOW3HXOh~EMo~A;{;y7D%S8KPU0oJj92g~Uc>8n18?Fjyp4D8E>7V+ypIn) z8;w*tAGHnV!k`cr3+W$WPd)vzQI4KGfBG6d@69h0y+fHst?{D#_w;8N_Z8CL=6hb?pFk@y Wm5Y=0e7g0qqn`iYird@Yum1pCAcHFa diff --git a/src/backend/scripts/reconcile_knowledge_document_projection.py b/src/backend/scripts/reconcile_knowledge_document_projection.py index 3d3726409..3b17d5650 100644 --- a/src/backend/scripts/reconcile_knowledge_document_projection.py +++ b/src/backend/scripts/reconcile_knowledge_document_projection.py @@ -1,6 +1,6 @@ """检查或重新调度单个 F059 文档入口投影。 -默认只读; 显式传入 ``--apply`` 才会向 knowledge Celery 队列调度任务: +默认只读; 显式传入 ``--apply`` 才会向默认 Celery 队列调度任务: PYTHONPATH=./ .venv/bin/python \ scripts/reconcile_knowledge_document_projection.py \ @@ -106,7 +106,7 @@ def main() -> int: "entry_id": int(args.entry_id), }, headers={"tenant_id": int(args.tenant_id)}, - queue="knowledge_celery", + queue="celery", ) snapshot["dispatch_status"] = "submitted" snapshot["task_id"] = str(task.id) diff --git a/src/backend/test/celery/test_knowledge_parse_queue_routing.py b/src/backend/test/celery/test_knowledge_parse_queue_routing.py new file mode 100644 index 000000000..272bbe6f7 --- /dev/null +++ b/src/backend/test/celery/test_knowledge_parse_queue_routing.py @@ -0,0 +1,173 @@ +import ast +from fnmatch import fnmatchcase +from pathlib import Path + +import pytest +import yaml + +from bisheng.core.config.celery_queues import ( + DEFAULT_CELERY_QUEUE, + KNOWLEDGE_PARSE_QUEUE, + KNOWLEDGE_PARSE_TASKS, + KNOWLEDGE_PDF_QUEUE, + PDF_ARTIFACT_TASK, + WORKFLOW_CELERY_QUEUE, + build_celery_task_routes, +) +from bisheng.core.config.settings import CeleryConf + +BACKEND_DIR = Path(__file__).resolve().parents[2] +PROJECT_DIR = Path(__file__).resolve().parents[4] + + +class _ConfigLoader(yaml.SafeLoader): + pass + + +_ConfigLoader.add_constructor("!env", lambda loader, node: loader.construct_scalar(node)) + + +def _resolve_queue(task_name: str, routes: dict) -> str | None: + route = routes.get(task_name) + if route is None: + route = next( + (candidate for pattern, candidate in routes.items() if "*" in pattern and fnmatchcase(task_name, pattern)), + None, + ) + if isinstance(route, str): + return route + return route.get("queue") if route else None + + +@pytest.mark.parametrize("task_name", sorted(KNOWLEDGE_PARSE_TASKS)) +def test_parse_task_whitelist_routes_to_knowledge_queue(task_name: str): + assert _resolve_queue(task_name, build_celery_task_routes({})) == KNOWLEDGE_PARSE_QUEUE + + +@pytest.mark.parametrize( + "task_name", + [ + "bisheng.worker.knowledge.file_worker.delete_knowledge_file_celery", + "bisheng.worker.knowledge.file_worker.copy_knowledge_file_celery", + "bisheng.worker.knowledge.file_worker.refresh_file_similarity_candidates_celery", + "bisheng.worker.knowledge.rebuild_knowledge_worker.rebuild_knowledge_celery", + "bisheng.worker.knowledge.qa.rebuild_qa_knowledge_celery", + "bisheng.worker.knowledge.file_migration.execute_knowledge_migration", + "bisheng.worker.knowledge.document_projection.process_document_projection", + "bisheng.worker.knowledge.portal_recommendation.rebuild_portal_recommendation_pool", + "bisheng.worker.org_sync.tasks.execute_org_sync", + "bisheng.worker.permission.department_transfer_cleanup.process_event", + "bisheng.worker.message.tasks.send_message", + "bisheng.worker.portal_course.tasks.scan_portal_course_media_cleanup", + ], +) +def test_non_parse_tasks_route_to_default_queue(task_name: str): + assert _resolve_queue(task_name, build_celery_task_routes({})) == DEFAULT_CELERY_QUEUE + + +@pytest.mark.parametrize( + ("task_name", "expected_queue"), + [ + (PDF_ARTIFACT_TASK, KNOWLEDGE_PDF_QUEUE), + ("bisheng.worker.workflow.tasks.run_workflow", WORKFLOW_CELERY_QUEUE), + ("bisheng.worker.approval.tasks.execute_approval_outbox", WORKFLOW_CELERY_QUEUE), + ], +) +def test_protected_non_default_routes_are_unchanged(task_name: str, expected_queue: str): + assert _resolve_queue(task_name, build_celery_task_routes({})) == expected_queue + + +def test_legacy_broad_knowledge_route_cannot_capture_non_parse_tasks(): + routes = build_celery_task_routes( + { + "bisheng.worker.knowledge.*": {"queue": "knowledge_celery"}, + "bisheng.worker.knowledge.document_projection.*": {"queue": "knowledge_celery"}, + "custom.task": {"queue": "custom_queue"}, + } + ) + + assert _resolve_queue(next(iter(KNOWLEDGE_PARSE_TASKS)), routes) == KNOWLEDGE_PARSE_QUEUE + assert _resolve_queue("bisheng.worker.knowledge.document_projection.reconcile", routes) == DEFAULT_CELERY_QUEUE + assert _resolve_queue("bisheng.worker.knowledge.qa.insert_qa_celery", routes) == DEFAULT_CELERY_QUEUE + assert _resolve_queue("custom.task", routes) == "custom_queue" + + +def test_celery_conf_uses_canonical_routes_for_default_and_legacy_config(): + default_config = CeleryConf() + legacy_config = CeleryConf( + task_routers={"bisheng.worker.knowledge.*": {"queue": "knowledge_celery"}}, + beat_schedule={}, + ) + + for config in (default_config, legacy_config): + for task_name in KNOWLEDGE_PARSE_TASKS: + assert _resolve_queue(task_name, config.task_routers) == KNOWLEDGE_PARSE_QUEUE + assert ( + _resolve_queue("bisheng.worker.knowledge.qa.insert_qa_celery", config.task_routers) == DEFAULT_CELERY_QUEUE + ) + + +@pytest.mark.parametrize( + "config_path", + [ + BACKEND_DIR / "bisheng/config_3002.yaml", + BACKEND_DIR / "bisheng/config_3003.yaml", + PROJECT_DIR / "docker/bisheng/config/config.yaml", + PROJECT_DIR / "docker/bisheng/config/config_dev.yaml", + ], +) +def test_runtime_yaml_declares_parse_only_knowledge_queue(config_path: Path): + config = yaml.load(config_path.read_text(encoding="utf-8"), Loader=_ConfigLoader) + routes = config["celery_task"]["task_routers"] + + assert routes["bisheng.worker.knowledge.*"] == {"queue": DEFAULT_CELERY_QUEUE} + for task_name in KNOWLEDGE_PARSE_TASKS: + assert routes[task_name] == {"queue": KNOWLEDGE_PARSE_QUEUE} + assert routes[PDF_ARTIFACT_TASK] == {"queue": KNOWLEDGE_PDF_QUEUE} + assert routes["bisheng.worker.workflow.*"] == {"queue": WORKFLOW_CELERY_QUEUE} + + +def test_worker_entrypoints_keep_default_queue_consumers_enabled(): + entrypoints = ( + BACKEND_DIR / "entrypoint.sh", + PROJECT_DIR / "docker/bisheng/entrypoint.sh", + ) + + for path in entrypoints: + source = path.read_text(encoding="utf-8") + assert "start_default" in source + assert "-Q celery" in source + + +def test_non_parse_production_dispatches_do_not_target_knowledge_queue(): + allowed_files = { + BACKEND_DIR / "scripts/enqueue_reparse_knowledge_space_files.py", + } + violations: list[str] = [] + + for root in (BACKEND_DIR / "bisheng", BACKEND_DIR / "scripts"): + for path in root.rglob("*.py"): + tree = ast.parse(path.read_text(encoding="utf-8"), filename=str(path)) + knowledge_queue_names = {"KNOWLEDGE_PARSE_QUEUE", "KNOWLEDGE_QUEUE"} | { + target.id + for node in ast.walk(tree) + if isinstance(node, (ast.Assign, ast.AnnAssign)) + for target in (node.targets if isinstance(node, ast.Assign) else [node.target]) + if isinstance(target, ast.Name) + and isinstance(node.value, ast.Constant) + and node.value.value == KNOWLEDGE_PARSE_QUEUE + } + for node in ast.walk(tree): + if not isinstance(node, ast.Call): + continue + queue_values = [keyword.value for keyword in node.keywords if keyword.arg == "queue"] + for value in queue_values: + targets_knowledge = ( + (isinstance(value, ast.Constant) and value.value == KNOWLEDGE_PARSE_QUEUE) + or (isinstance(value, ast.Name) and value.id in knowledge_queue_names) + or (isinstance(value, ast.Attribute) and value.attr == "KNOWLEDGE_PARSE_QUEUE") + ) + if targets_knowledge and path not in allowed_files: + violations.append(f"{path.relative_to(PROJECT_DIR)}:{node.lineno}") + + assert violations == [] diff --git a/src/backend/test/knowledge/test_favorite_notification_event.py b/src/backend/test/knowledge/test_favorite_notification_event.py index 6cdb349a3..8685df125 100644 --- a/src/backend/test/knowledge/test_favorite_notification_event.py +++ b/src/backend/test/knowledge/test_favorite_notification_event.py @@ -72,7 +72,7 @@ def test_enqueue_splits_event_batches(): second_payload = task.apply_async.call_args_list[1].kwargs["args"][0] assert len(first_payload) == 100 assert len(second_payload) == 1 - assert task.apply_async.call_args_list[0].kwargs["queue"] == "knowledge_celery" + assert task.apply_async.call_args_list[0].kwargs["queue"] == "celery" def test_enqueue_splits_large_delete_recipient_snapshot(): diff --git a/src/backend/test/knowledge/test_knowledge_document_permission_reconcile_worker.py b/src/backend/test/knowledge/test_knowledge_document_permission_reconcile_worker.py index 5943e167d..f8384745b 100644 --- a/src/backend/test/knowledge/test_knowledge_document_permission_reconcile_worker.py +++ b/src/backend/test/knowledge/test_knowledge_document_permission_reconcile_worker.py @@ -148,7 +148,7 @@ async def test_rollback_reconcile_dispatches_from_preparing_tombstone() -> None: "manager_file_id": 100, }, "headers": {"tenant_id": 7}, - "queue": projection_worker.KNOWLEDGE_QUEUE, + "queue": projection_worker.DEFAULT_QUEUE, } diff --git a/src/backend/test/knowledge/test_portal_hot_search_worker.py b/src/backend/test/knowledge/test_portal_hot_search_worker.py index 979c77941..d63ba8b0b 100644 --- a/src/backend/test/knowledge/test_portal_hot_search_worker.py +++ b/src/backend/test/knowledge/test_portal_hot_search_worker.py @@ -77,7 +77,7 @@ async def test_fanout_dispatches_per_tenant_with_headers(): assert count == 3 # default tenant + 5 + 7 tenant_ids = sorted(h["tenant_id"] for h, _q in dispatched) assert tenant_ids == [1, 5, 7] - assert all(q == "knowledge_celery" for _h, q in dispatched) + assert all(q == "celery" for _h, q in dispatched) @pytest.mark.asyncio diff --git a/src/backend/test/knowledge/test_portal_recommendation_worker.py b/src/backend/test/knowledge/test_portal_recommendation_worker.py index f0169518e..f4e60cacf 100644 --- a/src/backend/test/knowledge/test_portal_recommendation_worker.py +++ b/src/backend/test/knowledge/test_portal_recommendation_worker.py @@ -92,7 +92,7 @@ def test_tenant_fanout_always_uses_explicit_headers_and_knowledge_queue(): assert task.apply_async.call_count == 3 for tenant_id, call in zip((1, 5, 8), task.apply_async.call_args_list, strict=False): assert call.kwargs["headers"] == {"tenant_id": tenant_id} - assert call.kwargs["queue"] == "knowledge_celery" + assert call.kwargs["queue"] == "celery" assert call.kwargs["kwargs"] == {"mode": "daily"} diff --git a/src/backend/test/permission/test_department_transfer_cleanup_worker.py b/src/backend/test/permission/test_department_transfer_cleanup_worker.py index 8e1b9fff7..930fa1c8f 100644 --- a/src/backend/test/permission/test_department_transfer_cleanup_worker.py +++ b/src/backend/test/permission/test_department_transfer_cleanup_worker.py @@ -73,7 +73,7 @@ def test_tasks_are_late_acknowledged_and_permission_route_is_forced(): beat_schedule={}, ) - assert config.task_routers["bisheng.worker.permission.*"] == {"queue": "knowledge_celery"} + assert config.task_routers["bisheng.worker.permission.*"] == {"queue": "celery"} scan = config.beat_schedule["scan_department_transfer_permission_cleanup"] assert scan["task"].endswith(".scan_due_events") assert scan["schedule"] == 30.0 @@ -123,7 +123,7 @@ async def test_scan_recovers_preparing_event_and_dispatches_due_event(monkeypatc repository.activate_event.assert_awaited_once() worker_module.process_event.apply_async.assert_called_once_with( args=[51], - queue="knowledge_celery", + queue="celery", ) diff --git a/src/backend/test/shougang_portal_course/test_cleanup_service.py b/src/backend/test/shougang_portal_course/test_cleanup_service.py index 2ed2c7ab9..03486254c 100644 --- a/src/backend/test/shougang_portal_course/test_cleanup_service.py +++ b/src/backend/test/shougang_portal_course/test_cleanup_service.py @@ -111,7 +111,7 @@ def test_worker_and_beat_registration_are_stable(): config = CeleryConf() assert config.task_routers["bisheng.worker.portal_course.*"] == { - "queue": "knowledge_celery" + "queue": "celery" } assert config.beat_schedule["scan_portal_course_media_cleanup"]["task"] == ( "bisheng.worker.portal_course.tasks.scan_portal_course_media_cleanup" diff --git a/src/backend/test/telemetry_search/test_dashboard_enum_labels.py b/src/backend/test/telemetry_search/test_dashboard_enum_labels.py new file mode 100644 index 000000000..39cacf5af --- /dev/null +++ b/src/backend/test/telemetry_search/test_dashboard_enum_labels.py @@ -0,0 +1,348 @@ +from types import SimpleNamespace +from unittest.mock import AsyncMock + +import pytest + + +class _FakeDbSession: + async def __aenter__(self): + return SimpleNamespace() + + async def __aexit__(self, _exc_type, _exc, _traceback): + return False + + +async def _get_field_options( + monkeypatch, + *, + dataset_code, + field, + values, + label_field=None, + raw_labels=None, + keyword=None, +): + from bisheng.telemetry_search.domain.services import dashboard as module + + dimensions = [{"field": field, "field_type": "string"}] + if label_field: + dimensions.append({"field": label_field, "field_type": "string"}) + dataset = SimpleNamespace( + es_index_name=dataset_code, + schema_config={"dimensions": dimensions}, + ) + repository = SimpleNamespace(find_one=AsyncMock(return_value=dataset)) + buckets = [] + for index, value in enumerate(values): + bucket = {"key": value} + if label_field: + bucket["label_value"] = { + "buckets": [{"key": raw_labels[index]}], + } + buckets.append(bucket) + es_client = SimpleNamespace( + search=AsyncMock( + return_value={ + "aggregations": { + **( + { + "filter_wrapper": { + "total_count": {"value": len(values)}, + "enum_values": {"buckets": buckets}, + } + } + if keyword + else { + "total_count": {"value": len(values)}, + "enum_values": {"buckets": buckets}, + } + ) + } + } + ) + ) + + monkeypatch.setattr(module, "get_async_db_session", _FakeDbSession) + monkeypatch.setattr( + module, + "DashboardDatasetRepositoryImpl", + lambda _session: repository, + ) + monkeypatch.setattr( + module, + "get_es_connection", + AsyncMock(return_value=es_client), + ) + monkeypatch.setattr( + module.DashboardService, + "_get_realtime_scope_filters", + AsyncMock(return_value=[]), + ) + service = module.DashboardService.model_construct() + + result = await service.get_dataset_field_enums( + dataset_code=dataset_code, + field=field, + label_field=label_field, + keyword=keyword, + ) + return result, es_client.search.await_args.kwargs["body"] + + +@pytest.mark.parametrize( + ( + "dataset_code", + "field", + "label_field", + "values", + "raw_labels", + "expected_labels", + ), + [ + ( + "mid_app_increment", + "app_type", + None, + ["assistant", "workflow", "daily_chat", "unknown"], + None, + ["助手", "工作流", "日常对话", "未知"], + ), + ( + "mid_sessions_increment", + "source", + None, + ["platform", "api"], + None, + ["平台端", "API调用"], + ), + ( + "mid_sessions_increment", + "app_id", + "app_name", + ["daily_chat", "linsight", "42"], + ["daily_chat", "linsight", "自定义应用"], + ["日常对话", "Linsight", "自定义应用"], + ), + ( + "mid_session_run_dtl", + "app_id", + "app_name", + ["linsight"], + ["linsight"], + ["Linsight"], + ), + ( + "mid_tool_call_dtl", + "tool_type", + None, + [0, 1, 2], + None, + ["API工具", "内置工具", "MCP工具"], + ), + ( + "mid_tool_call_dtl", + "app_type", + None, + ["assistant"], + None, + ["助手"], + ), + ( + "mid_doc_parse_dtl", + "parse_type", + None, + ["local", "etl4lm", "un_etl4lm", "mineru", "paddle_ocr"], + None, + [ + "本地解析", + "ETL4LM解析", + "非ETL4LM解析", + "MinerU解析", + "PaddleOCR解析", + ], + ), + ( + "mid_doc_parse_dtl", + "status", + None, + ["success", "failed", "parse_failed"], + None, + ["成功", "失败", "解析失败"], + ), + ( + "mid_doc_parse_dtl", + "app_type", + None, + ["knowledge_base"], + None, + ["知识库"], + ), + ( + "mid_model_call_dtl", + "model_type", + None, + ["llm", "embedding", "rerank", "asr", "tts"], + None, + ["大语言模型", "嵌入模型", "重排模型", "语音识别", "语音合成"], + ), + ( + "mid_model_call_dtl", + "app_id", + "app_name", + ["evaluation"], + ["evaluation"], + ["模型评测"], + ), + ( + "mid_user_daily_participation", + "department_source", + None, + [ + "event_time", + "current_roster", + "current_roster_backfill", + "current_primary_backfill", + ], + None, + [ + "登录时所属主部门", + "当前在职名册", + "当前名册(历史回填)", + "当前主部门(历史登录回填)", + ], + ), + ], + ids=[ + "application-type", + "session-source", + "session-system-app", + "runtime-system-app", + "tool-type", + "tool-app-type", + "parse-type", + "parse-status", + "parse-app-type", + "model-type", + "model-system-app", + "participation-department-source", + ], +) +@pytest.mark.asyncio +async def test_dashboard_enum_options_keep_values_and_use_readable_labels( + monkeypatch, + dataset_code, + field, + label_field, + values, + raw_labels, + expected_labels, +): + result, _ = await _get_field_options( + monkeypatch, + dataset_code=dataset_code, + field=field, + values=values, + label_field=label_field, + raw_labels=raw_labels, + ) + + assert result["enums"] == values + assert result["options"] == [ + {"value": value, "label": label} for value, label in zip(values, expected_labels, strict=True) + ] + + +@pytest.mark.parametrize( + ("dataset_code", "field", "label_field", "keyword", "expected_values"), + [ + ("mid_app_increment", "app_type", None, "助手", ["assistant"]), + ( + "mid_sessions_increment", + "app_id", + "app_name", + "日常对话", + ["daily_chat"], + ), + ], +) +@pytest.mark.asyncio +async def test_dashboard_enum_search_matches_readable_labels( + monkeypatch, + dataset_code, + field, + label_field, + keyword, + expected_values, +): + result, body = await _get_field_options( + monkeypatch, + dataset_code=dataset_code, + field=field, + values=expected_values, + label_field=label_field, + raw_labels=expected_values if label_field else None, + keyword=keyword, + ) + + search_filter = body["aggs"]["filter_wrapper"]["filter"] + assert {"terms": {field: expected_values}} in search_filter["bool"]["should"] + assert result["enums"] == expected_values + + +def test_dashboard_dataset_seed_exposes_readable_tool_and_parse_dimensions(): + from bisheng.telemetry_search.domain.init_dataset import DASHBOARD_DATASET + + datasets = {dataset.dataset_code: dataset for dataset in DASHBOARD_DATASET} + tool_dimensions = { + dimension["field"]: dimension for dimension in datasets["mid_tool_call_dtl"].schema_config["dimensions"] + } + parse_dimensions = { + dimension["field"]: dimension for dimension in datasets["mid_doc_parse_dtl"].schema_config["dimensions"] + } + + assert tool_dimensions["app_type"]["name"] == "应用类型" + assert tool_dimensions["app_name"]["name"] == "应用名称" + assert parse_dimensions["app_type"]["name"] == "应用类型" + + +@pytest.mark.asyncio +async def test_dashboard_dataset_upgrade_refreshes_changed_dimension_schemas(): + from bisheng.telemetry_search.domain.init_dataset import ( + DASHBOARD_DATASET, + DASHBOARD_DATASET_REFRESH_CODES, + _upgrade_dashboard_datasets, + ) + + seeds = {dataset.dataset_code: dataset for dataset in DASHBOARD_DATASET} + existing = { + dataset_code: SimpleNamespace( + dataset_code=dataset_code, + dataset_name="旧名称", + es_index_name="旧索引", + description="旧描述", + is_commercial_only=False, + schema_config={}, + ) + for dataset_code in DASHBOARD_DATASET_REFRESH_CODES + } + repository = SimpleNamespace( + find_one=AsyncMock(side_effect=lambda dataset_code: existing.get(dataset_code)), + save=AsyncMock(), + update=AsyncMock(), + ) + + await _upgrade_dashboard_datasets(repository) + + assert repository.update.await_count == len(DASHBOARD_DATASET_REFRESH_CODES) + assert repository.save.await_count == 0 + for dataset_code, dataset in existing.items(): + assert dataset.schema_config == seeds[dataset_code].schema_config + + tool_dimensions = { + dimension["field"]: dimension for dimension in existing["mid_tool_call_dtl"].schema_config["dimensions"] + } + parse_dimensions = { + dimension["field"]: dimension for dimension in existing["mid_doc_parse_dtl"].schema_config["dimensions"] + } + assert tool_dimensions["app_name"]["name"] == "应用名称" + assert parse_dimensions["app_type"]["name"] == "应用类型" diff --git a/src/backend/test/test_realtime_dashboard.py b/src/backend/test/test_realtime_dashboard.py index d9cb1e263..345557df9 100644 --- a/src/backend/test/test_realtime_dashboard.py +++ b/src/backend/test/test_realtime_dashboard.py @@ -347,6 +347,181 @@ def test_file_space_level_options_group_section_libraries_under_team(): } +@pytest.mark.asyncio +async def test_date_field_enums_keep_raw_value_and_format_display_label( + monkeypatch, +): + from bisheng.telemetry_search.domain.services import dashboard as module + + timestamp_ms = 1785134769000 + dataset = SimpleNamespace( + es_index_name="mid_knowledge_space_content_stat", + schema_config={ + "dimensions": [ + { + "field": "timestamp", + "field_type": "date", + } + ] + }, + ) + repository = SimpleNamespace(find_one=AsyncMock(return_value=dataset)) + es_client = SimpleNamespace( + search=AsyncMock( + return_value={ + "aggregations": { + "total_count": {"value": 1}, + "enum_values": { + "buckets": [{"key": timestamp_ms}], + }, + } + } + ) + ) + + class FakeDbSession: + async def __aenter__(self): + return SimpleNamespace() + + async def __aexit__(self, _exc_type, _exc, _traceback): + return False + + monkeypatch.setattr(module, "get_async_db_session", FakeDbSession) + monkeypatch.setattr( + module, + "DashboardDatasetRepositoryImpl", + lambda _session: repository, + ) + monkeypatch.setattr( + module, + "get_es_connection", + AsyncMock(return_value=es_client), + ) + monkeypatch.setattr( + module.DashboardService, + "_get_realtime_scope_filters", + AsyncMock(return_value=[]), + ) + service = module.DashboardService.model_construct() + + result = await service.get_dataset_field_enums( + dataset_code="mid_knowledge_space_content_stat", + field="timestamp", + ) + + assert result["options"] == [ + { + "value": timestamp_ms, + "label": datetime.fromtimestamp(timestamp_ms / 1000).strftime( + "%Y-%m-%d %H:%M:%S" + ), + } + ] + + +@pytest.mark.parametrize( + ("field", "values", "expected_labels"), + [ + ( + "department_source", + ["event_time", "current_primary_backfill"], + ["提问时所属主部门", "当前主部门(历史回填)"], + ), + ( + "scene", + [ + "expert_question", + "smart_qa", + "document_qa", + "my_knowledge_document_qa", + ], + [ + "专家问答", + "智能问答", + "知识门户·文档问答", + "我的知识·文档问答", + ], + ), + ( + "source_app", + [ + "bisheng_my_knowledge", + "expert_qa", + "shougang_portal", + "unknown_app", + ], + ["毕昇·我的知识", "专家问答", "首钢知识门户", "unknown_app"], + ), + ], +) +@pytest.mark.asyncio +async def test_realtime_qa_field_enums_keep_codes_and_use_readable_labels( + monkeypatch, + field, + values, + expected_labels, +): + from bisheng.telemetry_search.domain.services import dashboard as module + + dataset = SimpleNamespace( + es_index_name="mid_realtime_qa_question_fact", + schema_config={ + "dimensions": [ + {"field": field, "field_type": "string"}, + ] + }, + ) + repository = SimpleNamespace(find_one=AsyncMock(return_value=dataset)) + es_client = SimpleNamespace( + search=AsyncMock( + return_value={ + "aggregations": { + "total_count": {"value": len(values)}, + "enum_values": { + "buckets": [{"key": value} for value in values], + }, + } + } + ) + ) + + class FakeDbSession: + async def __aenter__(self): + return SimpleNamespace() + + async def __aexit__(self, _exc_type, _exc, _traceback): + return False + + monkeypatch.setattr(module, "get_async_db_session", FakeDbSession) + monkeypatch.setattr( + module, + "DashboardDatasetRepositoryImpl", + lambda _session: repository, + ) + monkeypatch.setattr( + module, + "get_es_connection", + AsyncMock(return_value=es_client), + ) + monkeypatch.setattr( + module.DashboardService, + "_get_realtime_scope_filters", + AsyncMock(return_value=[]), + ) + service = module.DashboardService.model_construct() + + result = await service.get_dataset_field_enums( + dataset_code="mid_realtime_qa_question_fact", + field=field, + ) + + assert result["enums"] == values + assert result["options"] == [ + {"value": value, "label": label} + for value, label in zip(values, expected_labels, strict=True) + ] + + @pytest.mark.asyncio async def test_realtime_temporal_datasets_default_to_today_and_include_today(): from bisheng.telemetry_search.domain.schemas.component import ( diff --git a/src/frontend/platform/src/pages/Dashboard/components/charts/BaseChart.tsx b/src/frontend/platform/src/pages/Dashboard/components/charts/BaseChart.tsx index 1243d0a86..72c6b4cb8 100644 --- a/src/frontend/platform/src/pages/Dashboard/components/charts/BaseChart.tsx +++ b/src/frontend/platform/src/pages/Dashboard/components/charts/BaseChart.tsx @@ -199,11 +199,29 @@ const getPieChartOption = ( ) => { const { series } = data; const isDonut = chartType === 'donut'; + const metricNumberFormat = dataConfig && 'metrics' in dataConfig + ? dataConfig.metrics[0]?.numberFormat + : undefined; + const isPercentMetric = metricNumberFormat?.type === 'percent'; + + const formatMetricValue = (value: unknown) => { + if (!metricNumberFormat) return String(value); + const [formattedValue, unit] = unitConversion(value, dataConfig); + return `${formattedValue}${unit}`; + }; const tooltipFormatter = (params: any) => { - return `${params.name.replaceAll('\n', '
')}: ${params.value} (${params.percent}%)`; + const name = params.name.replaceAll('\n', '
'); + const formattedValue = formatMetricValue(params.value); + return isPercentMetric + ? `${name}: ${formattedValue}` + : `${name}: ${formattedValue} (${params.percent}%)`; }; + const dataLabelFormatter = isPercentMetric + ? (params: any) => `${params.name}: ${formatMetricValue(params.value)}` + : '{b}: {d}%'; + return { backgroundColor: styleConfig.bgColor, // title: buildTitleOption(styleConfig), @@ -221,7 +239,7 @@ const getPieChartOption = ( itemStyle: { borderRadius: 0, borderColor: '#fff', borderWidth: 2 }, label: { show: styleConfig.showDataLabel ?? true, - formatter: '{b}: {d}%', + formatter: dataLabelFormatter, fontSize: 10, color: "#666", textBorderWidth: 0, diff --git a/src/frontend/platform/src/test/pieChartTooltip.test.ts b/src/frontend/platform/src/test/pieChartTooltip.test.ts new file mode 100644 index 000000000..43c5b5004 --- /dev/null +++ b/src/frontend/platform/src/test/pieChartTooltip.test.ts @@ -0,0 +1,66 @@ +import { describe, expect, it } from "vitest" +import { generateChartOption } from "../pages/Dashboard/components/charts/BaseChart" +import { + ChartType, + ComponentStyleConfig, + DataConfig, +} from "../pages/Dashboard/types/dataConfig" + +describe("pie chart tooltip", () => { + it("formats a divide metric as a percentage without showing the slice share", () => { + const dataConfig: DataConfig = { + dimensions: [], + metrics: [ + { + fieldId: "participation_rate", + fieldName: "全员参与占比", + fieldCode: "participation_rate", + isVirtual: true, + sort: null, + numberFormat: { + type: "percent", + decimalPlaces: 2, + thousandSeparator: false, + }, + }, + ], + fieldOrder: [], + filters: [], + resultLimit: { limitType: "all" }, + isConfigured: true, + } + const option = generateChartOption({ + data: { + dimensions: [], + series: [ + { + name: "", + data: [ + { + name: "2026-08-03 ~ 2026-08-09", + value: 0.00023038820412394885, + }, + ], + }, + ], + }, + chartType: ChartType.Pie, + dataConfig, + styleConfig: {} as ComponentStyleConfig, + }) + + expect( + option.tooltip.formatter({ + name: "2026-08-03 ~ 2026-08-09", + value: 0.00023038820412394885, + percent: 100, + }), + ).toBe("2026-08-03 ~ 2026-08-09: 0.02%") + expect( + option.series[0].label.formatter({ + name: "2026-08-03 ~ 2026-08-09", + value: 0.00023038820412394885, + }), + ).toBe("2026-08-03 ~ 2026-08-09: 0.02%") + }) +})