From 32ceca31040bbb93486353f63c3461efbf817f9c Mon Sep 17 00:00:00 2001 From: Wenruli Date: Thu, 6 Aug 2026 13:33:07 +0800 Subject: [PATCH] feat: refine telemetry dashboard metrics --- docs/dashboard-dataset-metric-calculation.md | 4 +- features/v2.5.0-sg/release-contract.md | 4 +- .../bisheng/common/constants/telemetry.py | 1 + .../mid_table/knowledge_space_content.py | 15 +- .../telemetry_search/domain/init_dataset.py | 4 + .../domain/services/component.py | 58 +++--- .../domain/services/dashboard.py | 112 ++--------- .../bisheng/worker/telemetry/mid_table.py | 19 +- .../knowledge/test_portal_home_file_count.py | 1 + .../test_dashboard_enum_labels.py | 18 +- .../test_knowledge_space_content_telemetry.py | 53 +++++ src/backend/test/test_realtime_dashboard.py | 188 ++++++++---------- 12 files changed, 223 insertions(+), 254 deletions(-) diff --git a/docs/dashboard-dataset-metric-calculation.md b/docs/dashboard-dataset-metric-calculation.md index 9c6275fec..52f75ec1d 100644 --- a/docs/dashboard-dataset-metric-calculation.md +++ b/docs/dashboard-dataset-metric-calculation.md @@ -76,10 +76,10 @@ | 知识库文件存量表 | 总文件数 | `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)` | 当前有效文件快照按创建时间形成的累计数量 | 文件创建时间;支持年/月/周/日 | 当前每个文件只有一条快照 | +| 知识空间内容统计 | 总文件数 | `total_file_count` | `record_type=file`、`file_type=1`、`space_level` 为 `public`、`department`、`team`、`team_ks` 或 `personal`;排除“我的收藏” | 有时间维度:每桶 `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 | +| 知识空间内容统计 | 预览次数 | `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)` | 解析成功次数 ÷ 全部解析次数 | 解析事件时间;支持年/月/周/日 | 不额外去重 | diff --git a/features/v2.5.0-sg/release-contract.md b/features/v2.5.0-sg/release-contract.md index ecc8cbbde..9b7adda30 100644 --- a/features/v2.5.0-sg/release-contract.md +++ b/features/v2.5.0-sg/release-contract.md @@ -160,8 +160,8 @@ approval/share_link ────────────────┘ ### F058 不变式 -- **INV-SG-16**:统计投影和字段选项不得成为授权事实源;所有结果必须在后端按当前租户和当前用户权限重新裁剪。 -- **INV-SG-17**:文件总数只包含当前成功实体主版本文件;文件夹、失败、处理中、个人知识库和历史版本不得计入。 +- **INV-SG-16**:统计投影和字段选项不得成为资源授权事实源;看板资源访问继续由后端校验,统计数据范围由看板配置的租户、部门、知识空间等维度筛选决定,查询服务不得按当前用户隐式追加数据集范围条件。 +- **INV-SG-17**:文件总数只包含公共库、部门库、团队库、科室库和普通个人库中的当前成功实体主版本文件;文件夹、失败、处理中、“我的收藏”和历史版本不得计入。 - **INV-SG-18**:同一成功用户问题只生成一个问答事实;失败、重试和重新生成不得重复计数。 - **INV-SG-19**:参与人数按 Asia/Shanghai 自然日和 user_id 去重;Token 刷新不得写入登录事实。 - **INV-SG-20**:维度筛选只影响显式绑定的图表;旧查询组件和旧图表配置必须保持兼容。 diff --git a/src/backend/bisheng/common/constants/telemetry.py b/src/backend/bisheng/common/constants/telemetry.py index 1ed512eca..69e34dd20 100644 --- a/src/backend/bisheng/common/constants/telemetry.py +++ b/src/backend/bisheng/common/constants/telemetry.py @@ -5,4 +5,5 @@ KNOWLEDGE_SPACE_DASHBOARD_FILE_LEVELS = ( "department", "team", "team_ks", + "personal", ) diff --git a/src/backend/bisheng/telemetry/domain/mid_table/knowledge_space_content.py b/src/backend/bisheng/telemetry/domain/mid_table/knowledge_space_content.py index 0e730b868..913954d3d 100644 --- a/src/backend/bisheng/telemetry/domain/mid_table/knowledge_space_content.py +++ b/src/backend/bisheng/telemetry/domain/mid_table/knowledge_space_content.py @@ -687,7 +687,9 @@ return 0 viewer_user_name: str, occurred_at: datetime | None = None, ) -> None: - del space, viewer_user_id, viewer_user_name + if getattr(space, "is_favorite", False): + return + del viewer_user_id, viewer_user_name file_id = int(file_record.id) mid_table = cls(ensure_sync_index=False) try: @@ -757,19 +759,12 @@ return 0 raise RuntimeError(f"Failed to delete {len(real_errors)} knowledge space content file telemetry records.") return int(success or 0) - def delete_space_file_records_sync(self, space_ids: Iterable[int]) -> int: + def delete_space_records_sync(self, space_ids: Iterable[int]) -> int: ids = self._normalize_ids(space_ids) if not ids: return 0 result = self.delete_by_query_sync( - { - "bool": { - "filter": [ - {"term": {"record_type": "file"}}, - {"terms": {"space_id": ids}}, - ] - } - }, + {"terms": {"space_id": ids}}, refresh=False, ) return int(result.get("deleted", 0) or 0) diff --git a/src/backend/bisheng/telemetry_search/domain/init_dataset.py b/src/backend/bisheng/telemetry_search/domain/init_dataset.py index 6d4e6c865..acf4ec466 100644 --- a/src/backend/bisheng/telemetry_search/domain/init_dataset.py +++ b/src/backend/bisheng/telemetry_search/domain/init_dataset.py @@ -802,6 +802,10 @@ DASHBOARD_DATASET = [ is_virtual=True, filter=FilterExpression(bool_operator="must", filters=[ TermOp(field="record_type", value="preview_daily"), + TermsOp( + field="space_level", + value=list(KNOWLEDGE_SPACE_DASHBOARD_FILE_LEVELS), + ), ]), aggregations=[ AggregationExpression( diff --git a/src/backend/bisheng/telemetry_search/domain/services/component.py b/src/backend/bisheng/telemetry_search/domain/services/component.py index 531d5bcdb..b74f5d151 100644 --- a/src/backend/bisheng/telemetry_search/domain/services/component.py +++ b/src/backend/bisheng/telemetry_search/domain/services/component.py @@ -1,20 +1,43 @@ import copy from datetime import datetime, timedelta -from typing import List, Dict, Any, Union, Literal +from typing import Any, Dict, List, Literal, Union from loguru import logger from pydantic import BaseModel, Field -from bisheng.common.errcode.telemetry import QueryDatasetNotFoundError, QueryMetricNotFoundError, \ - QueryAggregationNotFoundError, QueryDimensionNotFoundError, QueryOperatorNotFoundError +from bisheng.common.errcode.telemetry import ( + QueryAggregationNotFoundError, + QueryDatasetNotFoundError, + QueryDimensionNotFoundError, + QueryMetricNotFoundError, + QueryOperatorNotFoundError, +) from bisheng.core.database import get_async_db_session -from .search_engine_service import SearchParameters, SearchEngineService -from ..models.dashboard_dataset import SchemaConfig, MetricConfig, DimensionConfig, FormulaEnum + +from ..models.dashboard_dataset import DimensionConfig, FormulaEnum, MetricConfig, SchemaConfig from ..repositories.implementations.dataset_repository_impl import DashboardDatasetRepositoryImpl -from ..schemas.component import ComponentDataConfig, TimeFilter, DataQueryResult, AggregationType, \ - DimensionField, DimensionQueryFilter, LogicType, OperatorType -from ..schemas.query_builder import AggregationExpression, AggsTypeEnum, FilterExpression, AtomFilter, TermOp, RangeOp, \ - RangeValue, MatchPhraseOp, TermsOp +from ..schemas.component import ( + AggregationType, + ComponentDataConfig, + DataQueryResult, + DimensionField, + DimensionQueryFilter, + LogicType, + OperatorType, + TimeFilter, +) +from ..schemas.query_builder import ( + AggregationExpression, + AggsTypeEnum, + AtomFilter, + FilterExpression, + MatchPhraseOp, + RangeOp, + RangeValue, + TermOp, + TermsOp, +) +from .search_engine_service import SearchEngineService, SearchParameters TIMESTAMP_FIELD = "timestamp" REALTIME_TEMPORAL_DATASETS = { @@ -31,10 +54,6 @@ class DataQueryService(BaseModel): default_factory=list, description="runtime dimension filters from linked filter components", ) - scope_filters: List[DimensionQueryFilter] = Field( - default_factory=list, - description="server-enforced tenant and administrative scope filters", - ) def _expand_dimension_filter_values( self, @@ -527,19 +546,6 @@ class DataQueryService(BaseModel): FilterExpression(bool_operator="must", filters=runtime_filters) ) - server_scope_filters = [ - TermsOp(field=scope_filter.field_id, value=scope_filter.values) - for scope_filter in self.scope_filters - if scope_filter.values - ] - if server_scope_filters: - filter_expressions.append( - FilterExpression( - bool_operator="must", - filters=server_scope_filters, - ) - ) - time_range = [] for one in all_time_filters: start_date, end_date = one.get_start_end_date( diff --git a/src/backend/bisheng/telemetry_search/domain/services/dashboard.py b/src/backend/bisheng/telemetry_search/domain/services/dashboard.py index 9ee29fe23..cc9578129 100644 --- a/src/backend/bisheng/telemetry_search/domain/services/dashboard.py +++ b/src/backend/bisheng/telemetry_search/domain/services/dashboard.py @@ -1,5 +1,5 @@ from datetime import datetime -from typing import List, Any, Sequence, Dict, ClassVar +from typing import Any, ClassVar, Dict, List, Sequence from fastapi import Request from pydantic import BaseModel, ConfigDict @@ -7,24 +7,20 @@ from sqlalchemy import Row, RowMapping from bisheng.api.services.audit_log import AuditLogService from bisheng.common.dependencies.user_deps import UserPayload -from bisheng.common.errcode.http_error import UnAuthorizedError, NotFoundError +from bisheng.common.errcode.http_error import NotFoundError, UnAuthorizedError from bisheng.common.errcode.telemetry import DashboardMaxError, DashBoardShareAuthError from bisheng.core.database import get_async_db_session -from bisheng.core.context.tenant import ( - DEFAULT_TENANT_ID, - get_current_tenant_id, -) from bisheng.core.search.elasticsearch.manager import get_es_connection -from bisheng.database.models.group_resource import GroupResourceDao, GroupResource, ResourceTypeEnum +from bisheng.database.models.group_resource import GroupResource, GroupResourceDao, ResourceTypeEnum from bisheng.database.models.role_access import AccessType, WebMenuResource from bisheng.user.domain.services.user import UserService from bisheng.utils import generate_uuid, get_request_ip -from ..models.dashboard import DashboardType, DashboardStatus, Dashboard, DashboardDefault, DashboardComponent + +from ..models.dashboard import Dashboard, DashboardComponent, DashboardDefault, DashboardStatus, DashboardType from ..models.dashboard_dao import DashboardDao from ..repositories.implementations.dataset_repository_impl import DashboardDatasetRepositoryImpl -from ..schemas.dashboard import DashboardRead, DashboardCreate -from ..schemas.component import DimensionQueryFilter -from ..services.component import TimeFilter, ComponentDataConfig, DataQueryService +from ..schemas.dashboard import DashboardCreate, DashboardRead +from ..services.component import ComponentDataConfig, DataQueryService, TimeFilter from ..utils import is_commercial @@ -34,7 +30,7 @@ class DashboardService(BaseModel): request: Request = None login_user: UserPayload = None - REALTIME_SCOPED_DATASETS: ClassVar[set[str]] = { + REALTIME_DATASETS: ClassVar[set[str]] = { "mid_knowledge_space_content_stat", "mid_realtime_qa_question_fact", "mid_user_daily_participation", @@ -43,6 +39,7 @@ class DashboardService(BaseModel): "public": "公共库", "department": "部门库", "team": "团队库(含科室库)", + "personal": "个人库", } APPLICATION_TYPE_LABELS: ClassVar[dict[str, str]] = { "workflow": "工作流", @@ -157,7 +154,7 @@ class DashboardService(BaseModel): components: Sequence[DashboardComponent], ) -> bool: return any( - component.dataset_code in cls.REALTIME_SCOPED_DATASETS + component.dataset_code in cls.REALTIME_DATASETS for component in components ) @@ -185,81 +182,6 @@ class DashboardService(BaseModel): ): raise UnAuthorizedError() - async def _get_realtime_scope_filters( - self, - dataset_code: str, - ) -> List[DimensionQueryFilter]: - if dataset_code not in self.REALTIME_SCOPED_DATASETS: - return [] - filters = [] - if dataset_code != "mid_knowledge_space_content_stat": - filters.append( - DimensionQueryFilter( - fieldId="tenant_id", - values=[get_current_tenant_id() or DEFAULT_TENANT_ID], - ) - ) - if self.login_user.is_admin(): - return filters - - from bisheng.database.models.department import DepartmentDao - - admin_departments = await DepartmentDao.aget_user_admin_departments( - self.login_user.user_id - ) - if not admin_departments: - raise UnAuthorizedError() - - if dataset_code == "mid_knowledge_space_content_stat": - from bisheng.common.models.space_channel_member import SpaceChannelMemberDao - from bisheng.permission.domain.services.permission_service import PermissionService - - accessible_ids = await PermissionService.list_accessible_ids( - user_id=self.login_user.user_id, - relation="can_manage", - object_type="knowledge_space", - login_user=self.login_user, - ) - managed_members = await SpaceChannelMemberDao.async_get_user_managed_members( - self.login_user.user_id - ) - space_ids = { - int(member.business_id) - for member in managed_members - if str(member.business_id).isdigit() - } - if accessible_ids is not None: - space_ids.update( - int(space_id) - for space_id in accessible_ids - if str(space_id).isdigit() - ) - filters.append( - DimensionQueryFilter( - fieldId="space_id", - values=sorted(space_ids) or ["__deny_all__"], - ) - ) - return filters - - department_ids = set() - for department in admin_departments: - department_ids.add(int(department.id)) - if department.path: - department_ids.update( - int(department_id) - for department_id in await DepartmentDao.aget_subtree_ids( - department.path - ) - ) - filters.append( - DimensionQueryFilter( - fieldId="primary_department_id", - values=sorted(department_ids) or ["__deny_all__"], - ) - ) - return filters - @classmethod async def get_simple_dashboards(cls, keyword: str = None, filter_ids: List[int] = None) -> List[Dashboard]: """ @@ -637,19 +559,17 @@ class DashboardService(BaseModel): if component is None: raise NotFoundError() if ( - component.dataset_code in self.REALTIME_SCOPED_DATASETS + component.dataset_code in self.REALTIME_DATASETS and not self.login_user.is_admin() and dashboard.status != DashboardStatus.PUBLISHED.value ): raise UnAuthorizedError() data_config = ComponentDataConfig(**component.data_config) - scope_filters = await self._get_realtime_scope_filters(component.dataset_code) res = await DataQueryService( dataset_code=component.dataset_code, data_config=data_config, time_filters=time_filters, dimension_filters=dimension_filters or [], - scope_filters=scope_filters, ).query_telemetry_data() return res @@ -701,7 +621,6 @@ class DashboardService(BaseModel): 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 es_client = await get_es_connection() @@ -789,14 +708,7 @@ class DashboardService(BaseModel): "size": 0, "aggs": aggs_body } - query_filters = [ - { - "terms": { - scope_filter.field_id: scope_filter.values, - } - } - for scope_filter in scope_filters - ] + query_filters = [] normalized_exact_values = [ value.strip() for value in (exact_values or "").split(",") diff --git a/src/backend/bisheng/worker/telemetry/mid_table.py b/src/backend/bisheng/worker/telemetry/mid_table.py index 1544deb37..605ba0a1d 100644 --- a/src/backend/bisheng/worker/telemetry/mid_table.py +++ b/src/backend/bisheng/worker/telemetry/mid_table.py @@ -386,6 +386,7 @@ def _get_success_space_file_rows(page: int, page_size: int): .join(Knowledge, KnowledgeFile.knowledge_id == Knowledge.id) .where( Knowledge.type == KnowledgeTypeEnum.SPACE.value, + Knowledge.is_favorite == False, # noqa: E712 KnowledgeFile.file_type == FileType.FILE.value, KnowledgeFile.status == KnowledgeFileStatus.SUCCESS.value, col(KnowledgeFile.deleted_at).is_(None), @@ -407,6 +408,7 @@ def _get_success_space_file_rows_by_space_id(space_id: int, page: int, page_size .where( Knowledge.id == space_id, Knowledge.type == KnowledgeTypeEnum.SPACE.value, + Knowledge.is_favorite == False, # noqa: E712 KnowledgeFile.file_type == FileType.FILE.value, KnowledgeFile.status == KnowledgeFileStatus.SUCCESS.value, col(KnowledgeFile.deleted_at).is_(None), @@ -438,6 +440,16 @@ def _get_knowledge_space_content_rows_by_file_ids(file_ids: List[int]): return session.exec(statement).all() +def _get_favorite_space_ids() -> list[int]: + statement = select(Knowledge.id).where( + Knowledge.type == KnowledgeTypeEnum.SPACE.value, + Knowledge.is_favorite == True, # noqa: E712 + ) + with bypass_tenant_filter(): + with get_sync_db_session() as session: + return [int(space_id) for space_id in session.exec(statement).all()] + + def _is_department_bound_space_scope(scope) -> bool: if scope is None: return False @@ -577,6 +589,7 @@ def _build_knowledge_space_content_records( def _is_file_content_stat_visible(file_record: KnowledgeFile, space: Knowledge) -> bool: return ( space.type == KnowledgeTypeEnum.SPACE.value + and not bool(getattr(space, "is_favorite", False)) and file_record.file_type == FileType.FILE.value and file_record.status == KnowledgeFileStatus.SUCCESS.value and getattr(file_record, "deleted_at", None) is None @@ -621,11 +634,15 @@ def rebuild_knowledge_space_content_file_projection(owner_token: str) -> dict[st if not KnowledgeSpaceContentStat.renew_lock_sync(owner_token): raise RuntimeError("Knowledge space content full projection owner lock lost before cleanup") deleted_count = mid_table.delete_stale_file_records_sync(sync_run_id) + deleted_favorite_count = mid_table.delete_space_records_sync( + _get_favorite_space_ids() + ) queue_status = KnowledgeSpaceContentStat.queue_status_sync() return { **queue_status, "synced": synced_count, "deleted_stale": deleted_count, + "deleted_favorite": deleted_favorite_count, "reclaimed_count": 0, "batch_duration_ms": KnowledgeSpaceContentStat._now_ms() - sync_started_ms, "projection_lag_ms": queue_status["oldest_pending_age_ms"], @@ -800,7 +817,7 @@ def sync_pending_knowledge_space_content_stat(): mid_table.insert_records_sync(records) space_synced_count += len(records) if space_synced_count == 0: - mid_table.delete_space_file_records_sync([space_id]) + mid_table.delete_space_records_sync([space_id]) if not KnowledgeSpaceContentStat.renew_lock_sync(owner_token): raise RuntimeError("Knowledge space content projection owner lock lost before space ack") if not KnowledgeSpaceContentStat.ack_claimed_sync(owner_token, [item.member]): diff --git a/src/backend/test/knowledge/test_portal_home_file_count.py b/src/backend/test/knowledge/test_portal_home_file_count.py index e1349b31f..fb6f56379 100644 --- a/src/backend/test/knowledge/test_portal_home_file_count.py +++ b/src/backend/test/knowledge/test_portal_home_file_count.py @@ -55,6 +55,7 @@ async def test_portal_file_count_uses_dashboard_es_metric(monkeypatch: pytest.Mo "department", "team", "team_ks", + "personal", ] } }, diff --git a/src/backend/test/telemetry_search/test_dashboard_enum_labels.py b/src/backend/test/telemetry_search/test_dashboard_enum_labels.py index 39cacf5af..fe2e72d12 100644 --- a/src/backend/test/telemetry_search/test_dashboard_enum_labels.py +++ b/src/backend/test/telemetry_search/test_dashboard_enum_labels.py @@ -73,11 +73,6 @@ async def _get_field_options( "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( @@ -252,6 +247,19 @@ async def test_dashboard_enum_options_keep_values_and_use_readable_labels( ] +@pytest.mark.asyncio +async def test_realtime_dataset_enum_query_has_no_implicit_scope_filter(monkeypatch): + result, body = await _get_field_options( + monkeypatch, + dataset_code="mid_realtime_qa_question_fact", + field="primary_department_id", + values=[9, 10, 11], + ) + + assert result["enums"] == [9, 10, 11] + assert "query" not in body + + @pytest.mark.parametrize( ("dataset_code", "field", "label_field", "keyword", "expected_values"), [ diff --git a/src/backend/test/test_knowledge_space_content_telemetry.py b/src/backend/test/test_knowledge_space_content_telemetry.py index 43f0a7837..fc93b320a 100644 --- a/src/backend/test/test_knowledge_space_content_telemetry.py +++ b/src/backend/test/test_knowledge_space_content_telemetry.py @@ -259,6 +259,28 @@ async def test_knowledge_space_content_log_preview_success_upserts_daily_counter }.intersection(call["upsert"]) +@pytest.mark.asyncio +async def test_favorite_space_preview_is_not_projected(monkeypatch): + from bisheng.telemetry.domain.mid_table import knowledge_space_content as module + + fake_client = _FakeAsyncIndexClient() + + async def fake_get_es_connection(): + return fake_client + + monkeypatch.setattr("bisheng.telemetry.domain.mid_table.base.get_es_connection", fake_get_es_connection) + + await module.KnowledgeSpaceContentStat.log_preview_success( + file_record=SimpleNamespace(id=11), + space=SimpleNamespace(id=3, name="我的收藏", is_favorite=True), + viewer_user_id=9, + viewer_user_name="查看人", + ) + + assert fake_client.get_calls == [] + assert fake_client.update_calls == [] + + @pytest.mark.asyncio async def test_knowledge_space_content_log_preview_success_does_not_retry_es_failure(monkeypatch): from bisheng.telemetry.domain.mid_table import knowledge_space_content as module @@ -384,6 +406,24 @@ def test_only_department_and_clinic_spaces_are_department_bound(level, expected) assert worker_module._is_department_bound_space_scope(SimpleNamespace(level=level)) is expected +@pytest.mark.parametrize( + ("is_favorite", "expected"), + [(False, True), (True, False)], +) +def test_content_projection_excludes_favorite_spaces(is_favorite, expected): + from bisheng.knowledge.domain.models.knowledge_file import FileType, KnowledgeFileStatus + + worker_module = _import_worker_mid_table() + file_record = SimpleNamespace( + file_type=FileType.FILE.value, + status=KnowledgeFileStatus.SUCCESS.value, + deleted_at=None, + ) + space = SimpleNamespace(type=3, is_favorite=is_favorite) + + assert worker_module._is_file_content_stat_visible(file_record, space) is expected + + def test_unbound_space_content_has_no_owning_department(): from bisheng.telemetry.domain.mid_table import knowledge_space_content as module @@ -434,6 +474,19 @@ def test_knowledge_space_content_delete_stale_file_records_uses_sync_run_id(monk assert call["body"]["query"]["bool"]["must_not"] == [{"term": {"sync_run_id": "run-1"}}] +def test_delete_space_records_removes_file_and_preview_rows(monkeypatch): + from bisheng.telemetry.domain.mid_table import knowledge_space_content as module + + fake_client = _FakeSyncIndexClient() + monkeypatch.setattr("bisheng.telemetry.domain.mid_table.base.get_es_connection_sync", lambda: fake_client) + + deleted = module.KnowledgeSpaceContentStat().delete_space_records_sync([3, 4]) + + assert deleted == 3 + query = fake_client.deleted_queries[0]["body"]["query"] + assert query == {"terms": {"space_id": [3, 4]}} + + def test_sync_pending_knowledge_space_content_stat_reloads_current_file_state(monkeypatch): from bisheng.knowledge.domain.models.knowledge_file import FileType, KnowledgeFileStatus diff --git a/src/backend/test/test_realtime_dashboard.py b/src/backend/test/test_realtime_dashboard.py index 345557df9..800fa9c9c 100644 --- a/src/backend/test/test_realtime_dashboard.py +++ b/src/backend/test/test_realtime_dashboard.py @@ -73,8 +73,26 @@ def test_realtime_dashboard_seed_contains_three_target_datasets(): assert knowledge_metrics["total_file_count"]["sum_type"] == "value_count" assert knowledge_metrics["new_file_count"]["aggregations"][0]["type"] == "value_count" assert knowledge_metrics["preview_count"]["filter"]["filters"] == [ - {"operator": "term", "field": "record_type", "value": "preview_daily"} + {"operator": "term", "field": "record_type", "value": "preview_daily"}, + { + "operator": "terms", + "field": "space_level", + "value": ["public", "department", "team", "team_ks", "personal"], + }, ] + for metric_name in ("total_file_count", "new_file_count", "contributor_count"): + level_filter = next( + item + for item in knowledge_metrics[metric_name]["filter"]["filters"] + if item["field"] == "space_level" + ) + assert level_filter["value"] == [ + "public", + "department", + "team", + "team_ks", + "personal", + ] assert knowledge_metrics["preview_count"]["aggregations"][0]["type"] == "sum" knowledge_dimensions = { dimension["field"]: dimension["name"] @@ -303,9 +321,6 @@ async def test_runtime_dimension_filters_are_combined_with_and(): DimensionQueryFilter(fieldId="space_level", values=["public", "team"]), DimensionQueryFilter(fieldId="business_domain_code", values=["steel"]), ], - scope_filters=[ - DimensionQueryFilter(fieldId="tenant_id", values=[2]), - ], ) filters, time_range = await service.convert_filters( { @@ -329,10 +344,66 @@ async def test_runtime_dimension_filters_are_combined_with_and(): "business_domain_code", ] assert dumped["filters"][0]["value"] == ["public", "team", "team_ks"] - scope = filters[1].model_dump() - assert scope["bool_operator"] == "must" - assert scope["filters"][0]["field"] == "tenant_id" - assert scope["filters"][0]["value"] == [2] + + +@pytest.mark.asyncio +async def test_realtime_component_query_does_not_pass_server_scope(monkeypatch): + from bisheng.telemetry_search.domain.models.dashboard import ( + Dashboard, + DashboardComponent, + DashboardStatus, + DashboardType, + ) + from bisheng.telemetry_search.domain.services import dashboard as module + + dashboard = Dashboard( + id=12, + title="实时统计", + status=DashboardStatus.PUBLISHED.value, + dashboard_type=DashboardType.PRESET_OSS.value, + user_id=7, + ) + component = DashboardComponent( + id="metric-1", + dashboard_id=12, + type="metric", + dataset_code="mid_realtime_qa_question_fact", + ) + captured_kwargs = {} + + class FakeDataQueryService: + def __init__(self, **kwargs): + captured_kwargs.update(kwargs) + + async def query_telemetry_data(self): + return "query-result" + + monkeypatch.setattr( + module.DashboardDao, + "get_one", + AsyncMock(return_value=dashboard), + ) + monkeypatch.setattr( + module.DashboardDao, + "get_one_component", + AsyncMock(return_value=component), + ) + monkeypatch.setattr(module, "DataQueryService", FakeDataQueryService) + service = module.DashboardService.model_construct( + login_user=SimpleNamespace( + is_admin=lambda: True, + async_access_check=AsyncMock(return_value=True), + ) + ) + + result = await service.query_component_data( + dashboard_id=12, + component_id="metric-1", + ) + + assert result == "query-result" + assert "scope_filters" not in captured_kwargs + assert captured_kwargs["dimension_filters"] == [] def test_file_space_level_options_group_section_libraries_under_team(): @@ -344,6 +415,7 @@ def test_file_space_level_options_group_section_libraries_under_team(): "public": "公共库", "department": "部门库", "team": "团队库(含科室库)", + "personal": "个人库", } @@ -397,11 +469,6 @@ async def test_date_field_enums_keep_raw_value_and_format_display_label( "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( @@ -503,11 +570,6 @@ async def test_realtime_qa_field_enums_keep_codes_and_use_readable_labels( "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( @@ -551,96 +613,6 @@ async def test_realtime_temporal_datasets_default_to_today_and_include_today(): assert datetime.fromtimestamp(seven_day_range[1] / 1000).date() == today -@pytest.mark.asyncio -async def test_knowledge_space_admin_scope_has_no_tenant_filter(monkeypatch): - from bisheng.telemetry_search.domain.services import dashboard as module - - monkeypatch.setattr(module, "get_current_tenant_id", lambda: 7) - service = module.DashboardService.model_construct( - login_user=SimpleNamespace(is_admin=lambda: True) - ) - - filters = await service._get_realtime_scope_filters( - "mid_knowledge_space_content_stat" - ) - - assert filters == [] - - -@pytest.mark.asyncio -async def test_department_admin_file_scope_uses_manageable_spaces(monkeypatch): - from bisheng.common.models.space_channel_member import SpaceChannelMemberDao - from bisheng.database.models.department import DepartmentDao - from bisheng.permission.domain.services.permission_service import ( - PermissionService, - ) - from bisheng.telemetry_search.domain.services import dashboard as module - - monkeypatch.setattr(module, "get_current_tenant_id", lambda: 7) - monkeypatch.setattr( - DepartmentDao, - "aget_user_admin_departments", - AsyncMock(return_value=[SimpleNamespace(id=9, path="1/9")]), - ) - monkeypatch.setattr( - PermissionService, - "list_accessible_ids", - AsyncMock(return_value=["12", "invalid"]), - ) - monkeypatch.setattr( - SpaceChannelMemberDao, - "async_get_user_managed_members", - AsyncMock(return_value=[SimpleNamespace(business_id="13")]), - ) - service = module.DashboardService.model_construct( - login_user=SimpleNamespace( - user_id=22, - is_admin=lambda: False, - ) - ) - - filters = await service._get_realtime_scope_filters( - "mid_knowledge_space_content_stat" - ) - - assert [item.model_dump(by_alias=True) for item in filters] == [ - {"fieldId": "space_id", "values": [12, 13]}, - ] - - -@pytest.mark.asyncio -async def test_department_admin_qa_scope_includes_department_subtree(monkeypatch): - from bisheng.database.models.department import DepartmentDao - from bisheng.telemetry_search.domain.services import dashboard as module - - monkeypatch.setattr(module, "get_current_tenant_id", lambda: 7) - monkeypatch.setattr( - DepartmentDao, - "aget_user_admin_departments", - AsyncMock(return_value=[SimpleNamespace(id=9, path="1/9")]), - ) - monkeypatch.setattr( - DepartmentDao, - "aget_subtree_ids", - AsyncMock(return_value=[9, 10, 11]), - ) - service = module.DashboardService.model_construct( - login_user=SimpleNamespace( - user_id=22, - is_admin=lambda: False, - ) - ) - - filters = await service._get_realtime_scope_filters( - "mid_realtime_qa_question_fact" - ) - - assert [item.model_dump(by_alias=True) for item in filters] == [ - {"fieldId": "tenant_id", "values": [7]}, - {"fieldId": "primary_department_id", "values": [9, 10, 11]}, - ] - - @pytest.mark.asyncio async def test_department_admin_cannot_edit_realtime_dashboard(monkeypatch): from bisheng.telemetry_search.domain.services import dashboard as module