from __future__ import annotations import os import uuid from datetime import UTC, datetime, timedelta import pytest from sqlalchemy import text pytestmark = pytest.mark.integration def test_quality_issue_postgres_persists_lifecycle_and_recurrence( monkeypatch, ): platform_url = os.environ.get("TEST_DATABASE_URL") if not platform_url: pytest.skip("TEST_DATABASE_URL is required") monkeypatch.setenv("DATABASE_URL", platform_url) from app import create_app, db from app.core.data_research.quality_issue_repository import ( SqlAlchemyQualityIssueRepository, ) from app.core.data_research.quality_issues import QualityIssueService from app.models.data_research import ( DeviceAsset, DeviceQualityIssue, DeviceQualityProfile, DeviceQualityProfileVersion, DeviceQualityRun, DeviceQualityViolationSample, IngestionSource, ) app = create_app() app.config.update(TESTING=True) suffix = uuid.uuid4().hex[:10] actor_uid = str(uuid.uuid4()) assignee_uid = str(uuid.uuid4()) reviewer_uid = str(uuid.uuid4()) source_uid = str(uuid.uuid4()) asset_uid = str(uuid.uuid4()) profile_uid = str(uuid.uuid4()) version_uid = str(uuid.uuid4()) run_uid = str(uuid.uuid4()) verify_run_uid = str(uuid.uuid4()) violation_uid = str(uuid.uuid4()) recurrence_violation_uid = str(uuid.uuid4()) issue_uids = [] now = datetime.now(UTC) try: with app.app_context(): for user_uid, username in ( (actor_uid, f"wp08-actor-{suffix}"), (assignee_uid, f"wp08-assignee-{suffix}"), (reviewer_uid, f"wp08-reviewer-{suffix}"), ): db.session.execute( text( """ INSERT INTO public.users ( id, username, display_name, password_hash, status ) VALUES ( CAST(:id AS uuid), :username, :username, 'integration-test', 'active' ) """ ), {"id": user_uid, "username": username}, ) db.session.add( IngestionSource( uid=source_uid, source_type="database", name=f"WP08 source {suffix}", config={"database_type": "postgresql"}, permission_scope={}, status="active", created_by=actor_uid, ) ) db.session.add( DeviceAsset( uid=asset_uid, asset_type="device", name=f"WP08 device {suffix}", status="active", current_version=1, content_hash="b" * 64, attributes={}, created_by=actor_uid, updated_by=actor_uid, ) ) db.session.add( DeviceQualityProfile( uid=profile_uid, code=f"WP08_{suffix}", name="WP08 integration profile", created_by=actor_uid, ) ) db.session.flush() db.session.add( DeviceQualityProfileVersion( uid=version_uid, profile_uid=profile_uid, version=1, status="draft", rules=[], content_hash="c" * 64, created_by=actor_uid, ) ) db.session.flush() for current_run_uid in (run_uid, verify_run_uid): db.session.add( DeviceQualityRun( uid=current_run_uid, policy_version_uid=version_uid, policy_hash="c" * 64, source_uid=source_uid, status="success", total_assets=1, total_violations=1, score=80, created_by=actor_uid, ) ) db.session.flush() for current_uid, current_run_uid in ( (violation_uid, run_uid), (recurrence_violation_uid, verify_run_uid), ): db.session.add( DeviceQualityViolationSample( uid=current_uid, run_uid=current_run_uid, rule_code="asset_context_complete", severity="error", asset_uid=asset_uid, field_name="location", source_uid=source_uid, message="设备位置缺失", evidence={ "asset_uid": asset_uid, "missing_fields": ["location"], }, created_at=now, expires_at=now + timedelta(days=30), ) ) db.session.commit() repository = SqlAlchemyQualityIssueRepository(db.session) service = QualityIssueService( repository, review_authorizer=lambda actor: ( None if actor == reviewer_uid else pytest.fail("unexpected reviewer") ), commit=db.session.commit, rollback=db.session.rollback, ) baseline_statistics = service.statistics() empty_records, empty_total = service.issues( status="closed", assignee_uid=assignee_uid, overdue_only=False, page=1, page_size=20, ) assert empty_records == [] assert empty_total == 0 imported = service.import_violations( violation_uids=[violation_uid], priority="high", due_at=now - timedelta(minutes=1), actor_uid=actor_uid, ) issue = imported.records[0] issue_uids.append(issue.uid) deduplicated = service.import_violations( violation_uids=[violation_uid], priority="critical", due_at=None, actor_uid=actor_uid, ) assert deduplicated.existing_count == 1 assigned = service.assign( issue.uid, assignee_uid=assignee_uid, due_at=issue.due_at, expected_version=1, actor_uid=actor_uid, ) started = service.start( issue.uid, expected_version=assigned.current_version, actor_uid=assignee_uid, ) submitted = service.submit( issue.uid, summary="设备位置已补录并复核来源", evidence_refs=[ { "type": "device_asset", "uid": asset_uid, "version": 2, } ], expected_version=started.current_version, actor_uid=assignee_uid, ) closed = service.review( issue.uid, verification_result="passed", note="复核通过", verification_run_uid=verify_run_uid, expected_version=submitted.current_version, actor_uid=reviewer_uid, ) assert closed.status == "closed" recurrence = service.import_violations( violation_uids=[recurrence_violation_uid], priority="high", due_at=None, actor_uid=actor_uid, ).records[0] issue_uids.append(recurrence.uid) assert recurrence.occurrence_number == 2 statistics = service.statistics() assert statistics["total"] == baseline_statistics["total"] + 2 assert statistics["recurrent"] == baseline_statistics["recurrent"] + 1 assert ( statistics["recurrent_issues"] == baseline_statistics["recurrent_issues"] + 1 ) assert statistics["recurrence_rate"] == round( statistics["recurrent_issues"] / statistics["total"], 6, ) assert [item.action for item in service.timeline(issue.uid)] == [ "created", "assigned", "started", "submitted", "closed", ] loaded, remediation = service.get(issue.uid) assert loaded.evidence["missing_fields"] == ["location"] assert remediation.review_status == "passed" assert remediation.verification_run_uid == verify_run_uid persisted = db.session.get(DeviceQualityIssue, issue.uid) assert persisted.current_version == 5 assert persisted.source_violation_uid == violation_uid finally: with app.app_context(): if issue_uids: db.session.query(DeviceQualityIssue).filter( DeviceQualityIssue.uid.in_(issue_uids) ).delete(synchronize_session=False) db.session.query(DeviceQualityViolationSample).filter( DeviceQualityViolationSample.uid.in_( [violation_uid, recurrence_violation_uid] ) ).delete(synchronize_session=False) db.session.query(DeviceQualityRun).filter( DeviceQualityRun.uid.in_([run_uid, verify_run_uid]) ).delete(synchronize_session=False) db.session.query(DeviceQualityProfileVersion).filter_by( uid=version_uid ).delete(synchronize_session=False) db.session.query(DeviceQualityProfile).filter_by( uid=profile_uid ).delete(synchronize_session=False) db.session.query(DeviceAsset).filter_by(uid=asset_uid).delete( synchronize_session=False ) db.session.query(IngestionSource).filter_by( uid=source_uid ).delete(synchronize_session=False) db.session.execute( text( "DELETE FROM public.users " "WHERE id = ANY(CAST(:ids AS uuid[]))" ), {"ids": [actor_uid, assignee_uid, reviewer_uid]}, ) db.session.commit()