DocWeave / backend /app /workflow /nodes /validate.py
shak3008's picture
fix: resolve workflow linking on resume, update proposals auto-commit, and validation baseline rules
90c285f
Raw
History Blame Contribute Delete
4.07 kB
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