mirror of
https://github.com/dataelement/bisheng.git
synced 2026-09-24 23:19:52 +08:00
feat: update dashboard data pipeline
This commit is contained in:
@@ -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
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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 核验运行时数据。
|
||||
@@ -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 |
|
||||
|
||||
---
|
||||
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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
|
||||
@@ -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",
|
||||
|
||||
@@ -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:
|
||||
# 通知是明确的尽力流程;入队失败不能回滚已完成的文档操作。
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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,
|
||||
)
|
||||
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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({
|
||||
|
||||
+1
-1
@@ -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")
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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"
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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,
|
||||
)
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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(
|
||||
|
||||
Binary file not shown.
@@ -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)
|
||||
|
||||
@@ -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 == []
|
||||
@@ -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():
|
||||
|
||||
@@ -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,
|
||||
}
|
||||
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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"}
|
||||
|
||||
|
||||
|
||||
@@ -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",
|
||||
)
|
||||
|
||||
|
||||
|
||||
@@ -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"
|
||||
|
||||
@@ -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"] == "应用类型"
|
||||
@@ -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 (
|
||||
|
||||
@@ -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', '<br/>')}: ${params.value} (${params.percent}%)`;
|
||||
const name = params.name.replaceAll('\n', '<br/>');
|
||||
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,
|
||||
|
||||
@@ -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%")
|
||||
})
|
||||
})
|
||||
Reference in New Issue
Block a user