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, } # Persist validation results to checkpoint for frontend access 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