| from __future__ import annotations |
|
|
| from sqlalchemy.orm import Session |
| from langgraph.graph import END, START, StateGraph |
| from langgraph.types import interrupt |
|
|
| from app.llm.client import LLMClient |
| from app.repositories.knowledge_repository import KnowledgeRepository |
| from app.services.chunking_service import ChunkingService |
| from app.services.knowledge_link_service import KnowledgeLinkService |
| from app.services.embedding_service import EmbeddingService |
|
|
| from app.workflow.nodes.chunk import chunk |
| from app.workflow.nodes.classify import classify |
| from app.workflow.nodes.complete import complete |
| from app.workflow.nodes.extract import extract |
| from app.workflow.nodes.knowledge import knowledge |
| from app.workflow.nodes.link import link |
| from app.workflow.nodes.reconcile import reconcile |
| from app.workflow.nodes.embed import embed |
| from app.workflow.nodes.validate import validate |
| from app.workflow.nodes.decision import decide |
|
|
| from app.workflow.router import route_after_decision |
| from app.workflow.state import WorkflowState |
|
|
|
|
| def build_workflow( |
| db: Session, |
| checkpointer=None, |
| interrupt_before=None, |
| ): |
| """ |
| Build the DocWeave document intelligence workflow. |
| |
| The human review branch uses LangGraph's durable interrupt() |
| mechanism so that the workflow state is checkpointed and can |
| later be resumed with Command(resume=...). |
| """ |
|
|
| llm_client = LLMClient() |
| knowledge_repository = KnowledgeRepository(db) |
| chunking_service = ChunkingService(db) |
| knowledge_link_service = KnowledgeLinkService(db) |
| embedding_service = EmbeddingService() |
|
|
| graph = StateGraph(WorkflowState) |
|
|
| |
| |
| |
|
|
| graph.add_node( |
| "extract", |
| extract, |
| ) |
|
|
| graph.add_node( |
| "chunk", |
| lambda state: chunk( |
| state, |
| chunking_service, |
| ), |
| ) |
|
|
| graph.add_node( |
| "embedding", |
| lambda state: embed( |
| state, |
| db, |
| embedding_service, |
| ), |
| ) |
|
|
| graph.add_node( |
| "classify", |
| classify, |
| ) |
|
|
| graph.add_node( |
| "knowledge", |
| lambda state: knowledge( |
| state, |
| llm_client, |
| knowledge_repository, |
| ), |
| ) |
|
|
| graph.add_node( |
| "reconciliation", |
| lambda state: reconcile( |
| state, |
| db, |
| ), |
| ) |
|
|
| graph.add_node( |
| "validation", |
| lambda state: validate( |
| state, |
| db, |
| ), |
| ) |
|
|
| graph.add_node( |
| "decision", |
| decide, |
| ) |
|
|
| graph.add_node( |
| "link", |
| lambda state: link( |
| state, |
| knowledge_link_service, |
| ), |
| ) |
|
|
| |
| |
| |
|
|
| def human_review( |
| state: WorkflowState, |
| ) -> WorkflowState: |
| """ |
| Durable human approval gate. |
| |
| LangGraph checkpoints the current state when interrupt() |
| is reached. The workflow can later continue from this |
| checkpoint using Command(resume=...). |
| """ |
|
|
| state.current_node = "HUMAN_REVIEW" |
|
|
| interrupt( |
| { |
| "type": "human_review", |
| "workflow_run_id": str( |
| state.workflow_run_id |
| ), |
| "document_version_id": str( |
| state.document_version_id |
| ), |
| "message": ( |
| "Human review is required before " |
| "the workflow can continue." |
| ), |
| } |
| ) |
|
|
| return state |
|
|
| graph.add_node( |
| "human_review", |
| human_review, |
| ) |
|
|
| graph.add_node( |
| "complete", |
| complete, |
| ) |
|
|
| |
| |
| |
|
|
| graph.add_edge( |
| START, |
| "extract", |
| ) |
|
|
| graph.add_edge( |
| "extract", |
| "chunk", |
| ) |
|
|
| graph.add_edge( |
| "chunk", |
| "embedding", |
| ) |
|
|
| graph.add_edge( |
| "embedding", |
| "classify", |
| ) |
|
|
| graph.add_edge( |
| "classify", |
| "knowledge", |
| ) |
|
|
| graph.add_edge( |
| "knowledge", |
| "reconciliation", |
| ) |
|
|
| graph.add_edge( |
| "reconciliation", |
| "validation", |
| ) |
|
|
| graph.add_edge( |
| "validation", |
| "decision", |
| ) |
|
|
| |
| |
| |
|
|
| graph.add_conditional_edges( |
| "decision", |
| route_after_decision, |
| { |
| "link": "link", |
| "human_review": "human_review", |
| }, |
| ) |
|
|
| |
| |
| |
|
|
| graph.add_edge( |
| "link", |
| "complete", |
| ) |
|
|
| |
| graph.add_edge( |
| "human_review", |
| "link", |
| ) |
|
|
| graph.add_edge( |
| "complete", |
| END, |
| ) |
|
|
| |
| |
| |
|
|
| return graph.compile( |
| checkpointer=checkpointer, |
| interrupt_before=interrupt_before, |
| ) |