| from __future__ import annotations |
|
|
| from app.agents.validation import RuleValidationAgent |
| from app.models.proposal import Proposal, ProposalType, ProposalStatus |
| from app.repositories.knowledge_repository import KnowledgeRepository |
| from app.repositories.rule_repository import RuleRepository |
| from app.workflow.state import WorkflowState |
|
|
|
|
| def validate( |
| state: WorkflowState, |
| db, |
| ) -> WorkflowState: |
| """ |
| Validate proposals generated by reconciliation against |
| enabled workspace rules. |
| """ |
|
|
| state.current_node = "RULE_VALIDATION" |
|
|
| from app.core.tracing.run_tracker import get_or_create_tracker |
| tracker = get_or_create_tracker(str(state.workflow_run_id)) |
| tracker.start_stage("validation") |
|
|
| from app.workflow.progress import report_progress |
| report_progress(state.workflow_run_id, "VALIDATION") |
|
|
| knowledge_repository = KnowledgeRepository(db) |
| rule_repository = RuleRepository(db) |
|
|
| new_items = knowledge_repository.list_by_document_version( |
| state.document_version_id |
| ) |
|
|
| new_item_ids = {str(item.id) for item in new_items} |
|
|
| if not new_item_ids: |
| error_msg = state.metadata.get("extraction_error") or "No knowledge items were extracted from this document." |
| has_err = bool(state.metadata.get("extraction_error")) |
| state.validation_results = [ |
| { |
| "status": "FAIL" if has_err else "WARNING", |
| "severity": "HIGH", |
| "rule_id": None, |
| "rule_name": "Extraction Verification", |
| "rule_type": None, |
| "operator": None, |
| "message": error_msg, |
| "proposal_ids": [], |
| } |
| ] |
| state.metadata["validation_summary"] = { |
| "status": "FAIL" if has_err else "WARNING", |
| "rules_evaluated": 1, |
| "proposals_evaluated": 0, |
| "failures": 1 if has_err else 0, |
| "warnings": 0 if has_err else 1, |
| } |
| tracker.end_stage("validation") |
| return state |
|
|
| proposals = ( |
| db.query(Proposal) |
| .filter( |
| Proposal.workspace_id == state.workspace_id, |
| Proposal.status == ProposalStatus.PENDING, |
| ) |
| .all() |
| ) |
|
|
| relevant_proposals = [] |
|
|
| for proposal in proposals: |
| changes = proposal.proposed_changes or {} |
|
|
| if proposal.proposal_type == ProposalType.CREATE: |
| if proposal.knowledge_item_id is not None: |
| if str(proposal.knowledge_item_id) in new_item_ids: |
| relevant_proposals.append(proposal) |
|
|
| elif proposal.proposal_type == ProposalType.UPDATE: |
| source_id = changes.get("source_knowledge_item_id") |
|
|
| if source_id and str(source_id) in new_item_ids: |
| relevant_proposals.append(proposal) |
|
|
| rules = rule_repository.list_enabled(state.workspace_id) |
|
|
| agent = RuleValidationAgent() |
|
|
| results = agent.validate( |
| rules=rules, |
| proposals=relevant_proposals, |
| ) |
|
|
| state.validation_results = results |
|
|
| failures = sum( |
| 1 for result in results if result["status"] == "FAIL" |
| ) |
|
|
| warnings = sum( |
| 1 for result in results if result["status"] == "WARNING" |
| ) |
|
|
| state.metadata["validation_summary"] = { |
| "status": ( |
| "FAIL" |
| if failures |
| else "WARNING" |
| if warnings |
| else "PASS" |
| ), |
| "rules_evaluated": len(rules), |
| "proposals_evaluated": len(relevant_proposals), |
| "failures": failures, |
| "warnings": warnings, |
| } |
|
|
| |
| from app.models.workflow_checkpoint import WorkflowCheckpoint |
| checkpoint = WorkflowCheckpoint( |
| workflow_run_id=state.workflow_run_id, |
| agent_name="VALIDATION", |
| state={"validation_results": results}, |
| message=f"Validation: {failures} fail, {warnings} warning", |
| ) |
| db.add(checkpoint) |
| db.commit() |
|
|
| tracker.end_stage("validation") |
| return state |
|
|