| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498 |
- from __future__ import annotations
- from dataclasses import replace
- from datetime import datetime
- import pytest
- class FakeDevelopmentService:
- def __init__(self, record):
- self.record = record
- self.actions = []
- def create_job(self, payload, actor_uid):
- self.actions.append(("create", payload, actor_uid))
- return replace(self.record, actor_uid=actor_uid), True
- def list_jobs(self, filters=None):
- self.actions.append(("list", filters or {}))
- return [self.record]
- def get_job(self, uid):
- self.actions.append(("get", uid))
- return self.record
- def retry(self, uid):
- self.actions.append(("retry", uid))
- return replace(self.record, status="queued", last_error=None)
- def cancel(self, uid):
- self.actions.append(("cancel", uid))
- return replace(self.record, status="cancelled")
- class FakeSourceRegistrar:
- def __init__(self):
- self.actions = []
- def ensure(self, source_uid, actor_uid):
- self.actions.append((source_uid, actor_uid))
- return object(), True
- class FakeCatalogExecutor:
- def __init__(self, record):
- self.record = record
- self.actions = []
- def execute(self, uid):
- self.actions.append(uid)
- return replace(
- self.record,
- status="awaiting_review",
- attempt_count=1,
- statistics={
- "snapshot_uid": "snapshot-1",
- "asset_count": 2,
- "field_count": 8,
- "evidence_count": 8,
- },
- )
- class FakeCatalogSnapshotRepository:
- def list(self, job_uid):
- from app.core.data_research.catalog.models import CatalogSnapshotRecord
- return [
- CatalogSnapshotRecord(
- uid="snapshot-1",
- job_uid=job_uid,
- source_uid="00000000-0000-0000-0000-000000000001",
- attempt=1,
- database_type="postgresql",
- content_hash="b" * 64,
- snapshot={
- "data_source_uid": "00000000-0000-0000-0000-000000000001",
- "database_type": "postgresql",
- "assets": [],
- },
- evidence_count=8,
- )
- ]
- class FakeEvidenceService:
- def list_for_job(self, job_uid):
- return [
- {
- "uid": "evidence-1",
- "job_uid": job_uid,
- "locator": {
- "kind": "database.column",
- "schema": "asset",
- "table": "equipment",
- "column": "equipment_code",
- },
- "excerpt": "password=[redacted]",
- "confidence": 1.0,
- }
- ]
- class FakeDeviceAssetService:
- def __init__(self):
- from app.core.data_research.device_assets import (
- DeviceAssetDetail,
- DeviceAssetImportItem,
- DeviceAssetImportResult,
- DeviceAssetMappingRecord,
- DeviceAssetRecord,
- DeviceAssetVersionRecord,
- )
- now = datetime(2026, 7, 29, 11, 0)
- self.asset = DeviceAssetRecord(
- uid="00000000-0000-7000-8000-000000000301",
- asset_type="device",
- name="一号循环泵",
- status="active",
- current_version=2,
- content_hash="c" * 64,
- location="动力车间",
- organization="设备动力部",
- responsible_person="张工",
- attributes={"model": "P-100"},
- created_by="editor-1",
- updated_by="editor-1",
- created_at=now,
- updated_at=now,
- )
- self.mapping = DeviceAssetMappingRecord(
- uid="00000000-0000-7000-8000-000000000302",
- asset_uid=self.asset.uid,
- source_uid="00000000-0000-0000-0000-000000000101",
- source_entity="asset.equipment",
- asset_type="device",
- source_code="EQ-001",
- source_updated_at=now,
- first_seen_at=now,
- last_seen_at=now,
- )
- self.version = DeviceAssetVersionRecord(
- uid="00000000-0000-7000-8000-000000000303",
- asset_uid=self.asset.uid,
- version=2,
- content_hash=self.asset.content_hash,
- snapshot={
- "asset_type": "device",
- "source_code": "EQ-001",
- "name": "一号循环泵",
- "status": "active",
- "location": "动力车间",
- "organization": "设备动力部",
- "responsible_person": "张工",
- "attributes": {"model": "P-100"},
- },
- source_mapping_uid=self.mapping.uid,
- actor_uid="editor-1",
- created_at=now,
- )
- self.detail = DeviceAssetDetail(
- asset=self.asset,
- mappings=(self.mapping,),
- )
- self.import_result = DeviceAssetImportResult(
- items=(
- DeviceAssetImportItem(
- action="updated",
- asset=self.asset,
- mapping=self.mapping,
- ),
- ),
- created_count=0,
- updated_count=1,
- unchanged_count=0,
- )
- self.actions = []
- def import_records(self, payload, actor_uid):
- self.actions.append(("import", payload, actor_uid))
- return self.import_result
- def search(self, filters, *, page, page_size):
- self.actions.append(("search", filters, page, page_size))
- return [self.detail], 1
- def get(self, asset_uid):
- self.actions.append(("get", asset_uid))
- return self.detail
- def versions(self, asset_uid):
- self.actions.append(("versions", asset_uid))
- return [self.version]
- @pytest.fixture()
- def development_client(monkeypatch):
- from flask import request
- from app import create_app
- from app.api.data_development import routes
- from app.core.data_research.models import IngestionJobRecord
- from app.core.system import permissions
- record = IngestionJobRecord(
- uid="00000000-0000-0000-0000-000000000010",
- source_uid="00000000-0000-0000-0000-000000000001",
- artifact_uid="00000000-0000-0000-0000-000000000002",
- job_type="file_extract",
- parser_version="csv-v1",
- idempotency_key="a" * 64,
- actor_uid="editor-1",
- parameters={"schema": "public"},
- )
- service = FakeDevelopmentService(record)
- registrar = FakeSourceRegistrar()
- executor = FakeCatalogExecutor(record)
- asset_service = FakeDeviceAssetService()
- def identity():
- header = request.headers.get("Authorization", "")
- if header == "Bearer viewer":
- return {"id": "viewer-1", "roles": ["viewer"]}
- if header == "Bearer editor":
- return {"id": "editor-1", "roles": ["editor"]}
- if header == "Bearer other-editor":
- return {"id": "editor-2", "roles": ["editor"]}
- if header == "Bearer admin":
- return {"id": "admin-1", "roles": ["admin"]}
- return None
- monkeypatch.setattr(permissions, "authenticate_request", identity)
- monkeypatch.setattr(routes, "get_ingestion_service", lambda: service)
- monkeypatch.setattr(
- routes,
- "get_database_source_registration_service",
- lambda: registrar,
- )
- monkeypatch.setattr(
- routes,
- "get_catalog_ingestion_executor",
- lambda: executor,
- )
- monkeypatch.setattr(
- routes,
- "get_catalog_snapshot_repository",
- lambda: FakeCatalogSnapshotRepository(),
- )
- monkeypatch.setattr(
- routes,
- "get_evidence_service",
- lambda: FakeEvidenceService(),
- )
- monkeypatch.setattr(
- routes,
- "get_device_asset_service",
- lambda: asset_service,
- raising=False,
- )
- app = create_app()
- app.config.update(TESTING=True)
- service.registrar = registrar
- service.executor = executor
- service.asset_service = asset_service
- return app.test_client(), service
- def job_payload():
- return {
- "source_uid": "00000000-0000-0000-0000-000000000001",
- "artifact_uid": "00000000-0000-0000-0000-000000000002",
- "job_type": "file_extract",
- "parser_version": "csv-v1",
- "parameters": {"schema": "public"},
- "password": "must-not-return",
- }
- def test_ingestion_api_requires_authentication_and_run_permission(development_client):
- client, _service = development_client
- assert client.post("/api/development/v1/ingestion-jobs", json=job_payload()).status_code == 401
- assert client.post(
- "/api/development/v1/ingestion-jobs",
- json=job_payload(),
- headers={"Authorization": "Bearer viewer"},
- ).status_code == 403
- def test_editor_creates_and_lists_secret_free_jobs(development_client):
- client, service = development_client
- created = client.post(
- "/api/development/v1/ingestion-jobs",
- json=job_payload(),
- headers={"Authorization": "Bearer editor"},
- )
- listed = client.get(
- "/api/development/v1/ingestion-jobs?status=created",
- headers={"Authorization": "Bearer editor"},
- )
- assert created.status_code == 201
- assert created.get_json()["data"]["uid"].endswith("0010")
- assert "must-not-return" not in created.get_data(as_text=True)
- assert listed.status_code == 200
- assert listed.get_json()["data"]["total"] == 1
- assert service.actions[-1] == ("list", {"status": "created"})
- def test_retry_is_admin_only_and_cancel_is_owner_or_admin(development_client):
- client, _service = development_client
- uid = "00000000-0000-0000-0000-000000000010"
- editor_retry = client.post(
- f"/api/development/v1/ingestion-jobs/{uid}/retry",
- headers={"Authorization": "Bearer editor"},
- )
- admin_retry = client.post(
- f"/api/development/v1/ingestion-jobs/{uid}/retry",
- headers={"Authorization": "Bearer admin"},
- )
- other_cancel = client.post(
- f"/api/development/v1/ingestion-jobs/{uid}/cancel",
- headers={"Authorization": "Bearer other-editor"},
- )
- owner_cancel = client.post(
- f"/api/development/v1/ingestion-jobs/{uid}/cancel",
- headers={"Authorization": "Bearer editor"},
- )
- assert editor_retry.status_code == 403
- assert admin_retry.status_code == 200
- assert other_cancel.status_code == 403
- assert owner_cancel.status_code == 200
- def test_catalog_job_registers_source_then_executes_with_report(
- development_client,
- ):
- client, service = development_client
- payload = {
- **job_payload(),
- "artifact_uid": None,
- "job_type": "catalog_collect",
- "parser_version": "catalog-v1",
- }
- created = client.post(
- "/api/development/v1/ingestion-jobs",
- json=payload,
- headers={"Authorization": "Bearer editor"},
- )
- executed = client.post(
- "/api/development/v1/ingestion-jobs/"
- "00000000-0000-0000-0000-000000000010/execute",
- headers={"Authorization": "Bearer editor"},
- )
- assert created.status_code == 201
- assert service.registrar.actions == [
- ("00000000-0000-0000-0000-000000000001", "editor-1")
- ]
- assert executed.status_code == 200
- assert executed.get_json()["data"]["status"] == "awaiting_review"
- assert executed.get_json()["data"]["attempt_count"] == 1
- assert service.executor.actions == [
- "00000000-0000-0000-0000-000000000010"
- ]
- def test_catalog_snapshots_and_evidence_are_readable_without_secrets(
- development_client,
- ):
- client, _service = development_client
- uid = "00000000-0000-0000-0000-000000000010"
- snapshots = client.get(
- f"/api/development/v1/ingestion-jobs/{uid}/catalog-snapshots",
- headers={"Authorization": "Bearer viewer"},
- )
- evidence = client.get(
- f"/api/development/v1/ingestion-jobs/{uid}/evidence",
- headers={"Authorization": "Bearer viewer"},
- )
- assert snapshots.status_code == 200
- assert snapshots.get_json()["data"]["records"][0]["attempt"] == 1
- assert evidence.status_code == 200
- text = evidence.get_data(as_text=True)
- assert "clear-secret" not in text
- assert "equipment_code" in text
- def test_device_asset_catalog_is_readable_and_import_requires_edit_permission(
- development_client,
- ):
- client, service = development_client
- endpoint = "/api/development/v1/device-assets"
- payload = {
- "source_uid": "00000000-0000-0000-0000-000000000101",
- "source_entity": "asset.equipment",
- "records": [
- {
- "asset_type": "device",
- "source_code": "EQ-001",
- "name": "一号循环泵",
- "attributes": {"model": "P-100"},
- }
- ],
- }
- viewed = client.get(
- f"{endpoint}?keyword=循环泵&asset_type=device&status=active"
- "&source_uid=00000000-0000-0000-0000-000000000101"
- "&page=2&page_size=10",
- headers={"Authorization": "Bearer viewer"},
- )
- denied = client.post(
- f"{endpoint}/import",
- json=payload,
- headers={"Authorization": "Bearer viewer"},
- )
- imported = client.post(
- f"{endpoint}/import",
- json=payload,
- headers={"Authorization": "Bearer editor"},
- )
- assert viewed.status_code == 200
- data = viewed.get_json()["data"]
- assert data["page"] == 2
- assert data["page_size"] == 10
- assert data["total"] == 1
- assert data["records"][0]["uid"].endswith("0301")
- assert data["records"][0]["source_mappings"][0]["source_code"] == "EQ-001"
- assert service.asset_service.actions[0] == (
- "search",
- {
- "keyword": "循环泵",
- "asset_type": "device",
- "status": "active",
- "source_uid": "00000000-0000-0000-0000-000000000101",
- },
- 2,
- 10,
- )
- assert denied.status_code == 403
- assert imported.status_code == 200
- assert imported.get_json()["data"]["updated_count"] == 1
- assert service.asset_service.actions[-1][2] == "editor-1"
- def test_device_asset_detail_and_versions_return_traceability_without_secrets(
- development_client,
- ):
- client, _service = development_client
- uid = "00000000-0000-7000-8000-000000000301"
- detail = client.get(
- f"/api/development/v1/device-assets/{uid}",
- headers={"Authorization": "Bearer viewer"},
- )
- versions = client.get(
- f"/api/development/v1/device-assets/{uid}/versions",
- headers={"Authorization": "Bearer viewer"},
- )
- assert detail.status_code == 200
- assert versions.status_code == 200
- assert detail.get_json()["data"]["current_version"] == 2
- assert detail.get_json()["data"]["source_mappings"][0][
- "source_entity"
- ] == "asset.equipment"
- assert versions.get_json()["data"]["records"][0]["version"] == 2
- text = detail.get_data(as_text=True) + versions.get_data(as_text=True)
- assert "password" not in text.lower()
- assert "credential" not in text.lower()
- def test_device_asset_catalog_rejects_unbounded_or_malformed_pagination(
- development_client,
- ):
- client, _service = development_client
- endpoint = "/api/development/v1/device-assets"
- malformed = client.get(
- f"{endpoint}?page=not-a-number",
- headers={"Authorization": "Bearer viewer"},
- )
- oversized = client.get(
- f"{endpoint}?page_size=101",
- headers={"Authorization": "Bearer viewer"},
- )
- assert malformed.status_code == 422
- assert oversized.status_code == 422
- assert malformed.get_json()["error"]["code"] == "DEVICE_ASSET_INVALID"
|