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