feat: refine telemetry dashboard metrics

This commit is contained in:
Wenruli
2026-08-06 13:33:07 +08:00
parent 457badb445
commit 32ceca3104
12 changed files with 223 additions and 254 deletions
+2 -2
View File
@@ -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)` | 解析成功次数 ÷ 全部解析次数 | 解析事件时间;支持年/月/周/日 | 不额外去重 |
+2 -2
View File
@@ -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**:维度筛选只影响显式绑定的图表;旧查询组件和旧图表配置必须保持兼容。
@@ -5,4 +5,5 @@ KNOWLEDGE_SPACE_DASHBOARD_FILE_LEVELS = (
"department",
"team",
"team_ks",
"personal",
)
@@ -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)
@@ -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(
@@ -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(
@@ -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(",")
@@ -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]):
@@ -55,6 +55,7 @@ async def test_portal_file_count_uses_dashboard_es_metric(monkeypatch: pytest.Mo
"department",
"team",
"team_ks",
"personal",
]
}
},
@@ -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"),
[
@@ -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
+80 -108
View File
@@ -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