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) # ------------------------------------------------------------------ # Nodes # ------------------------------------------------------------------ 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, ), ) # ------------------------------------------------------------------ # Durable human review gate # ------------------------------------------------------------------ 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, ) # ------------------------------------------------------------------ # Workflow edges # ------------------------------------------------------------------ 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", ) # ------------------------------------------------------------------ # Decision routing # ------------------------------------------------------------------ graph.add_conditional_edges( "decision", route_after_decision, { "link": "link", "human_review": "human_review", }, ) # ------------------------------------------------------------------ # Normal completion path # ------------------------------------------------------------------ graph.add_edge( "link", "complete", ) # After human approval/resume, route through link before completion. graph.add_edge( "human_review", "link", ) graph.add_edge( "complete", END, ) # ------------------------------------------------------------------ # Compile # ------------------------------------------------------------------ return graph.compile( checkpointer=checkpointer, interrupt_before=interrupt_before, )