shak3008 commited on
Commit
0868c75
·
1 Parent(s): d1644cd

feat: add human review approval gate and workflow resume

Browse files
backend/alembic/versions/add_waiting_for_review.py ADDED
@@ -0,0 +1,27 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+ """add waiting for review workflow status
2
+
3
+ Revision ID: add_waiting_for_review
4
+ Revises: 585aa57edb27
5
+ """
6
+
7
+ from alembic import op
8
+
9
+
10
+ revision = "add_waiting_for_review"
11
+ down_revision = "585aa57edb27"
12
+ branch_labels = None
13
+ depends_on = None
14
+
15
+
16
+ def upgrade() -> None:
17
+ op.execute(
18
+ "ALTER TYPE workflow_status "
19
+ "ADD VALUE IF NOT EXISTS 'WAITING_FOR_REVIEW'"
20
+ )
21
+
22
+
23
+ def downgrade() -> None:
24
+ # PostgreSQL does not support removing a value directly
25
+ # from an enum type. A downgrade would require recreating
26
+ # the enum and migrating the column.
27
+ pass
backend/app/api/workflows.py ADDED
@@ -0,0 +1,60 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+ from uuid import UUID
2
+
3
+ from fastapi import APIRouter, Depends, HTTPException
4
+ from sqlalchemy.orm import Session
5
+
6
+ from app.core.dependencies import get_current_user
7
+ from app.database.session import get_db
8
+ from app.models.user import User
9
+ from app.models.workflow_run import WorkflowRun
10
+ from app.models.workspace import Workspace
11
+
12
+ router = APIRouter(
13
+ prefix="/workflows",
14
+ tags=["Workflows"],
15
+ )
16
+
17
+
18
+ @router.get("/{workflow_id}")
19
+ def get_workflow(
20
+ workflow_id: UUID,
21
+ current_user: User = Depends(get_current_user),
22
+ db: Session = Depends(get_db),
23
+ ):
24
+ workflow = (
25
+ db.query(WorkflowRun)
26
+ .filter(WorkflowRun.id == workflow_id)
27
+ .first()
28
+ )
29
+
30
+ if workflow is None:
31
+ raise HTTPException(
32
+ status_code=404,
33
+ detail="Workflow not found.",
34
+ )
35
+
36
+ workspace = (
37
+ db.query(Workspace)
38
+ .filter(
39
+ Workspace.id == workflow.workspace_id,
40
+ Workspace.created_by == current_user.id,
41
+ )
42
+ .first()
43
+ )
44
+
45
+ if workspace is None:
46
+ raise HTTPException(
47
+ status_code=403,
48
+ detail="You do not have access to this workflow.",
49
+ )
50
+
51
+ return {
52
+ "id": str(workflow.id),
53
+ "workspace_id": str(workflow.workspace_id),
54
+ "document_version_id": str(
55
+ workflow.document_version_id
56
+ ),
57
+ "status": workflow.status.value,
58
+ "started_at": workflow.started_at,
59
+ "completed_at": workflow.completed_at,
60
+ }
backend/app/main.py CHANGED
@@ -1,14 +1,17 @@
1
- from fastapi import FastAPI
2
  from fastapi.middleware.cors import CORSMiddleware
 
 
3
  from app.api.workspaces import router as workspace_router
4
- from app.api import auth, documents
5
- from app.api import proposals
6
  app = FastAPI(
7
  title="DocWeave",
8
  description="Agentic Document Intelligence Platform",
9
  version="0.1.0",
10
  )
11
 
 
12
  app.add_middleware(
13
  CORSMiddleware,
14
  allow_origins=["*"],
@@ -17,10 +20,26 @@ app.add_middleware(
17
  allow_headers=["*"],
18
  )
19
 
20
- app.include_router(auth.router, prefix="/auth", tags=["auth"])
21
- app.include_router(documents.router, prefix="/documents", tags=["documents"])
 
 
 
 
 
 
 
 
 
 
 
22
  app.include_router(workspace_router)
 
23
  app.include_router(proposals.router)
 
 
 
 
24
  @app.get("/")
25
  def root():
26
  return {
 
1
+ from fastapi import FastAPI
2
  from fastapi.middleware.cors import CORSMiddleware
3
+
4
+ from app.api import auth, documents, proposals, workflows
5
  from app.api.workspaces import router as workspace_router
6
+
7
+
8
  app = FastAPI(
9
  title="DocWeave",
10
  description="Agentic Document Intelligence Platform",
11
  version="0.1.0",
12
  )
13
 
14
+
15
  app.add_middleware(
16
  CORSMiddleware,
17
  allow_origins=["*"],
 
20
  allow_headers=["*"],
21
  )
22
 
23
+
24
+ app.include_router(
25
+ auth.router,
26
+ prefix="/auth",
27
+ tags=["auth"],
28
+ )
29
+
30
+ app.include_router(
31
+ documents.router,
32
+ prefix="/documents",
33
+ tags=["documents"],
34
+ )
35
+
36
  app.include_router(workspace_router)
37
+
38
  app.include_router(proposals.router)
39
+
40
+ app.include_router(workflows.router)
41
+
42
+
43
  @app.get("/")
44
  def root():
45
  return {
backend/app/models/workflow_run.py CHANGED
@@ -19,6 +19,7 @@ from app.database.database import Base
19
  class WorkflowStatus(str, enum.Enum):
20
  PENDING = "PENDING"
21
  RUNNING = "RUNNING"
 
22
  COMPLETED = "COMPLETED"
23
  FAILED = "FAILED"
24
  CANCELLED = "CANCELLED"
 
19
  class WorkflowStatus(str, enum.Enum):
20
  PENDING = "PENDING"
21
  RUNNING = "RUNNING"
22
+ WAITING_FOR_REVIEW = "WAITING_FOR_REVIEW"
23
  COMPLETED = "COMPLETED"
24
  FAILED = "FAILED"
25
  CANCELLED = "CANCELLED"
backend/app/repositories/workflow_repository.py CHANGED
@@ -18,6 +18,23 @@ class WorkflowRepository:
18
 
19
  return workflow
20
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
21
  def update_status(
22
  self,
23
  workflow: WorkflowRun,
 
18
 
19
  return workflow
20
 
21
+ def get(self, workflow_id):
22
+ return (
23
+ self.db.query(WorkflowRun)
24
+ .filter(WorkflowRun.id == workflow_id)
25
+ .first()
26
+ )
27
+
28
+ def get_by_document_version(self, document_version_id):
29
+ return (
30
+ self.db.query(WorkflowRun)
31
+ .filter(
32
+ WorkflowRun.document_version_id == document_version_id,
33
+ )
34
+ .order_by(WorkflowRun.started_at.desc())
35
+ .first()
36
+ )
37
+
38
  def update_status(
39
  self,
40
  workflow: WorkflowRun,
backend/app/services/document_service.py CHANGED
@@ -12,6 +12,7 @@ from app.storage.local_storage import LocalStorage
12
  from app.services.workflow_service import WorkflowService
13
  from app.workflow import WorkflowExecutor, WorkflowState
14
 
 
15
  class DocumentService:
16
  def __init__(self, db: Session):
17
  self.repository = DocumentRepository(db)
@@ -51,9 +52,14 @@ class DocumentService:
51
 
52
  self.repository.refresh(document)
53
  self.repository.refresh(version)
54
- #todo:
55
- # Move workflow orchestration out of DocumentService.
56
- # The upload service should only persist documents.
 
 
 
 
 
57
  workflow_service = WorkflowService(self.repository.db)
58
 
59
  workflow = workflow_service.start_workflow(
@@ -68,7 +74,12 @@ class DocumentService:
68
  document_path=version.storage_path,
69
  )
70
 
71
- executor = WorkflowExecutor(self.repository.db)
 
 
 
 
 
72
 
73
  try:
74
  executor.execute(
@@ -77,4 +88,5 @@ class DocumentService:
77
  )
78
  finally:
79
  executor.close()
 
80
  return document, version
 
12
  from app.services.workflow_service import WorkflowService
13
  from app.workflow import WorkflowExecutor, WorkflowState
14
 
15
+
16
  class DocumentService:
17
  def __init__(self, db: Session):
18
  self.repository = DocumentRepository(db)
 
52
 
53
  self.repository.refresh(document)
54
  self.repository.refresh(version)
55
+
56
+ # TODO:
57
+ # Move workflow orchestration out of DocumentService.
58
+ #
59
+ # The upload service should eventually only persist documents.
60
+ # Workflow execution should be handled by a dedicated workflow
61
+ # application/service layer.
62
+
63
  workflow_service = WorkflowService(self.repository.db)
64
 
65
  workflow = workflow_service.start_workflow(
 
74
  document_path=version.storage_path,
75
  )
76
 
77
+ # The workflow pauses before the human_review node when the
78
+ # decision agent determines that human review is required.
79
+ executor = WorkflowExecutor(
80
+ self.repository.db,
81
+ interrupt_before=["human_review"],
82
+ )
83
 
84
  try:
85
  executor.execute(
 
88
  )
89
  finally:
90
  executor.close()
91
+
92
  return document, version
backend/app/services/proposal_review_service.py CHANGED
@@ -10,9 +10,8 @@ from app.models.commit import Commit
10
  from app.models.knowledge_item import KnowledgeItem, KnowledgeStatus
11
  from app.models.proposal import Proposal, ProposalStatus, ProposalType
12
  from app.models.review import Review, ReviewDecision
13
- from app.models.user import User
14
  from app.models.workspace import Workspace
15
-
16
 
17
  class ProposalReviewService:
18
  """
@@ -50,6 +49,169 @@ class ProposalReviewService:
50
  .all()
51
  )
52
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
53
  def approve(
54
  self,
55
  proposal_id: uuid.UUID,
@@ -90,6 +252,8 @@ class ProposalReviewService:
90
  self.db.commit()
91
  self.db.refresh(commit)
92
 
 
 
93
  return commit
94
 
95
  except Exception:
@@ -125,8 +289,9 @@ class ProposalReviewService:
125
  self.db.commit()
126
  self.db.refresh(review)
127
 
128
- return review
129
 
 
130
  except Exception:
131
  self.db.rollback()
132
  raise
@@ -226,15 +391,16 @@ class ProposalReviewService:
226
 
227
  existing_item.status = KnowledgeStatus.ACTIVE
228
 
229
- # The newly extracted item was the evidence that caused
230
- # this proposal. It should not become a second active
231
- # canonical item.
232
  source_id = changes.get("source_knowledge_item_id")
233
 
234
  if source_id:
235
- source_item = self.db.query(KnowledgeItem).filter(
236
- KnowledgeItem.id == uuid.UUID(source_id)
237
- ).first()
 
 
 
 
238
 
239
  if source_item is not None:
240
  source_item.status = KnowledgeStatus.SUPERSEDED
@@ -308,9 +474,7 @@ class ProposalReviewService:
308
  if proposal.status != ProposalStatus.PENDING:
309
  raise HTTPException(
310
  status_code=409,
311
- detail=(
312
- "Proposal has already been reviewed."
313
- ),
314
  )
315
 
316
  def _verify_workspace_access(
@@ -330,7 +494,5 @@ class ProposalReviewService:
330
  if workspace is None:
331
  raise HTTPException(
332
  status_code=403,
333
- detail=(
334
- "You do not have access to this workspace."
335
- ),
336
  )
 
10
  from app.models.knowledge_item import KnowledgeItem, KnowledgeStatus
11
  from app.models.proposal import Proposal, ProposalStatus, ProposalType
12
  from app.models.review import Review, ReviewDecision
 
13
  from app.models.workspace import Workspace
14
+ from app.models.workflow_run import WorkflowStatus
15
 
16
  class ProposalReviewService:
17
  """
 
49
  .all()
50
  )
51
 
52
+ def get_pending_for_document_version(
53
+ self,
54
+ workspace_id: uuid.UUID,
55
+ document_version_id: uuid.UUID,
56
+ ) -> list[Proposal]:
57
+ """
58
+ Return pending proposals associated with a specific document version.
59
+
60
+ CREATE proposals normally point directly at the newly extracted
61
+ KnowledgeItem.
62
+
63
+ UPDATE proposals point at the existing canonical KnowledgeItem,
64
+ so their source_knowledge_item_id inside proposed_changes is also
65
+ checked.
66
+ """
67
+
68
+ knowledge_items = (
69
+ self.db.query(KnowledgeItem.id)
70
+ .filter(
71
+ KnowledgeItem.workspace_id == workspace_id,
72
+ KnowledgeItem.document_version_id == document_version_id,
73
+ )
74
+ .all()
75
+ )
76
+
77
+ knowledge_item_ids = {
78
+ str(item_id)
79
+ for (item_id,) in knowledge_items
80
+ }
81
+
82
+ if not knowledge_item_ids:
83
+ return []
84
+
85
+ proposals = (
86
+ self.db.query(Proposal)
87
+ .filter(
88
+ Proposal.workspace_id == workspace_id,
89
+ Proposal.status == ProposalStatus.PENDING,
90
+ )
91
+ .order_by(Proposal.created_at.asc())
92
+ .all()
93
+ )
94
+
95
+ relevant = []
96
+
97
+ for proposal in proposals:
98
+ # CREATE / DELETE / directly-associated proposals
99
+ if (
100
+ proposal.knowledge_item_id is not None
101
+ and str(proposal.knowledge_item_id)
102
+ in knowledge_item_ids
103
+ ):
104
+ relevant.append(proposal)
105
+ continue
106
+
107
+ # UPDATE proposals reference the newly extracted source item
108
+ # through proposed_changes.
109
+ changes = proposal.proposed_changes or {}
110
+
111
+ source_id = changes.get("source_knowledge_item_id")
112
+
113
+ if source_id and str(source_id) in knowledge_item_ids:
114
+ relevant.append(proposal)
115
+
116
+ return relevant
117
+
118
+ def _get_proposal_document_version_id(
119
+ self,
120
+ proposal: Proposal,
121
+ ) -> uuid.UUID | None:
122
+ """
123
+ Resolve the document version that generated this proposal.
124
+
125
+ CREATE/DELETE proposals normally point directly at the extracted
126
+ KnowledgeItem.
127
+
128
+ UPDATE proposals point at the existing canonical item, so the
129
+ newly extracted source KnowledgeItem is stored in
130
+ proposed_changes["source_knowledge_item_id"].
131
+ """
132
+
133
+ source_item_id = None
134
+
135
+ changes = proposal.proposed_changes or {}
136
+
137
+ if proposal.proposal_type == ProposalType.UPDATE:
138
+ source_item_id = changes.get(
139
+ "source_knowledge_item_id"
140
+ )
141
+
142
+ if source_item_id:
143
+ source_item = (
144
+ self.db.query(KnowledgeItem)
145
+ .filter(
146
+ KnowledgeItem.id == uuid.UUID(
147
+ str(source_item_id)
148
+ )
149
+ )
150
+ .first()
151
+ )
152
+
153
+ if source_item is not None:
154
+ return source_item.document_version_id
155
+
156
+ if proposal.knowledge_item_id is not None:
157
+ item = (
158
+ self.db.query(KnowledgeItem)
159
+ .filter(
160
+ KnowledgeItem.id == proposal.knowledge_item_id
161
+ )
162
+ .first()
163
+ )
164
+
165
+ if item is not None:
166
+ return item.document_version_id
167
+
168
+ return None
169
+
170
+ def _resume_workflow_if_ready(
171
+ self,
172
+ proposal: Proposal,
173
+ ) -> None:
174
+ document_version_id = (
175
+ self._get_proposal_document_version_id(
176
+ proposal
177
+ )
178
+ )
179
+
180
+ if document_version_id is None:
181
+ return
182
+
183
+ pending = self.get_pending_for_document_version(
184
+ proposal.workspace_id,
185
+ document_version_id,
186
+ )
187
+
188
+ if pending:
189
+ return
190
+
191
+ from app.repositories.workflow_repository import (
192
+ WorkflowRepository,
193
+ )
194
+ from app.workflow import WorkflowExecutor
195
+
196
+ workflow_repository = WorkflowRepository(self.db)
197
+
198
+ workflow = workflow_repository.get_by_document_version(
199
+ document_version_id
200
+ )
201
+
202
+ if workflow is None:
203
+ return
204
+
205
+ if workflow.status != WorkflowStatus.WAITING_FOR_REVIEW:
206
+ return
207
+
208
+ executor = WorkflowExecutor(self.db)
209
+
210
+ try:
211
+ executor.resume(workflow)
212
+ finally:
213
+ executor.close()
214
+
215
  def approve(
216
  self,
217
  proposal_id: uuid.UUID,
 
252
  self.db.commit()
253
  self.db.refresh(commit)
254
 
255
+ self._resume_workflow_if_ready(proposal)
256
+
257
  return commit
258
 
259
  except Exception:
 
289
  self.db.commit()
290
  self.db.refresh(review)
291
 
292
+ self._resume_workflow_if_ready(proposal)
293
 
294
+ return review
295
  except Exception:
296
  self.db.rollback()
297
  raise
 
391
 
392
  existing_item.status = KnowledgeStatus.ACTIVE
393
 
 
 
 
394
  source_id = changes.get("source_knowledge_item_id")
395
 
396
  if source_id:
397
+ source_item = (
398
+ self.db.query(KnowledgeItem)
399
+ .filter(
400
+ KnowledgeItem.id == uuid.UUID(source_id)
401
+ )
402
+ .first()
403
+ )
404
 
405
  if source_item is not None:
406
  source_item.status = KnowledgeStatus.SUPERSEDED
 
474
  if proposal.status != ProposalStatus.PENDING:
475
  raise HTTPException(
476
  status_code=409,
477
+ detail="Proposal has already been reviewed.",
 
 
478
  )
479
 
480
  def _verify_workspace_access(
 
494
  if workspace is None:
495
  raise HTTPException(
496
  status_code=403,
497
+ detail="You do not have access to this workspace.",
 
 
498
  )
backend/app/services/workflow_service.py CHANGED
@@ -25,6 +25,9 @@ class WorkflowService:
25
 
26
  return workflow
27
 
 
 
 
28
  def complete_workflow(self, workflow):
29
  workflow = self.repository.update_status(
30
  workflow,
@@ -34,4 +37,15 @@ class WorkflowService:
34
  self.repository.commit()
35
  self.repository.refresh(workflow)
36
 
 
 
 
 
 
 
 
 
 
 
 
37
  return workflow
 
25
 
26
  return workflow
27
 
28
+ def get_workflow(self, workflow_id):
29
+ return self.repository.get(workflow_id)
30
+
31
  def complete_workflow(self, workflow):
32
  workflow = self.repository.update_status(
33
  workflow,
 
37
  self.repository.commit()
38
  self.repository.refresh(workflow)
39
 
40
+ return workflow
41
+
42
+ def wait_for_review(self, workflow):
43
+ workflow = self.repository.update_status(
44
+ workflow,
45
+ WorkflowStatus.WAITING_FOR_REVIEW,
46
+ )
47
+
48
+ self.repository.commit()
49
+ self.repository.refresh(workflow)
50
+
51
  return workflow
backend/app/workflow/executor.py CHANGED
@@ -1,9 +1,10 @@
1
  from __future__ import annotations
2
 
3
  import psycopg
4
- from sqlalchemy.orm import Session
5
  from langgraph.checkpoint.postgres import PostgresSaver
6
  from langgraph.checkpoint.serde.jsonplus import JsonPlusSerializer
 
 
7
 
8
  from app.core.config import DATABASE_URL
9
  from app.services.workflow_service import WorkflowService
@@ -48,6 +49,27 @@ class WorkflowExecutor:
48
  interrupt_before=interrupt_before,
49
  )
50
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
51
  def execute(
52
  self,
53
  state: WorkflowState | None,
@@ -65,9 +87,43 @@ class WorkflowExecutor:
65
  config=config,
66
  )
67
 
68
- self.workflow_service.complete_workflow(workflow)
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
69
 
70
  return final_state
71
 
 
 
 
 
72
  def close(self):
73
  self.checkpoint_connection.close()
 
1
  from __future__ import annotations
2
 
3
  import psycopg
 
4
  from langgraph.checkpoint.postgres import PostgresSaver
5
  from langgraph.checkpoint.serde.jsonplus import JsonPlusSerializer
6
+ from langgraph.types import Command
7
+ from sqlalchemy.orm import Session
8
 
9
  from app.core.config import DATABASE_URL
10
  from app.services.workflow_service import WorkflowService
 
49
  interrupt_before=interrupt_before,
50
  )
51
 
52
+ # ------------------------------------------------------------------
53
+ # Helpers
54
+ # ------------------------------------------------------------------
55
+
56
+ @staticmethod
57
+ def _is_completed(state) -> bool:
58
+ """
59
+ LangGraph may return either a WorkflowState instance or
60
+ a dictionary depending on the configured state schema and
61
+ execution path.
62
+ """
63
+
64
+ if isinstance(state, dict):
65
+ return bool(state.get("completed", False))
66
+
67
+ return bool(getattr(state, "completed", False))
68
+
69
+ # ------------------------------------------------------------------
70
+ # Initial execution
71
+ # ------------------------------------------------------------------
72
+
73
  def execute(
74
  self,
75
  state: WorkflowState | None,
 
87
  config=config,
88
  )
89
 
90
+ if self._is_completed(final_state):
91
+ self.workflow_service.complete_workflow(workflow)
92
+ else:
93
+ self.workflow_service.wait_for_review(workflow)
94
+
95
+ return final_state
96
+
97
+ # ------------------------------------------------------------------
98
+ # Resume after human review
99
+ # ------------------------------------------------------------------
100
+
101
+ def resume(
102
+ self,
103
+ workflow,
104
+ ) -> WorkflowState:
105
+
106
+ config = {
107
+ "configurable": {
108
+ "thread_id": str(workflow.id),
109
+ }
110
+ }
111
+
112
+ final_state = self.graph.invoke(
113
+ Command(resume=True),
114
+ config=config,
115
+ )
116
+
117
+ if self._is_completed(final_state):
118
+ self.workflow_service.complete_workflow(workflow)
119
+ else:
120
+ self.workflow_service.wait_for_review(workflow)
121
 
122
  return final_state
123
 
124
+ # ------------------------------------------------------------------
125
+ # Cleanup
126
+ # ------------------------------------------------------------------
127
+
128
  def close(self):
129
  self.checkpoint_connection.close()
backend/app/workflow/graph.py CHANGED
@@ -2,6 +2,7 @@
2
 
3
  from sqlalchemy.orm import Session
4
  from langgraph.graph import END, START, StateGraph
 
5
 
6
  from app.llm.client import LLMClient
7
  from app.repositories.knowledge_repository import KnowledgeRepository
@@ -31,6 +32,10 @@ def build_workflow(
31
  ):
32
  """
33
  Build the DocWeave document intelligence workflow.
 
 
 
 
34
  """
35
 
36
  llm_client = LLMClient()
@@ -110,10 +115,39 @@ def build_workflow(
110
  ),
111
  )
112
 
113
- # Temporary branch target.
114
- # This will become the real durable human gate later.
115
- def human_review(state: WorkflowState) -> WorkflowState:
 
 
 
 
 
 
 
 
 
 
 
 
116
  state.current_node = "HUMAN_REVIEW"
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
117
  return state
118
 
119
  graph.add_node(
@@ -170,7 +204,10 @@ def build_workflow(
170
  "decision",
171
  )
172
 
173
- # Decision controls the workflow path.
 
 
 
174
  graph.add_conditional_edges(
175
  "decision",
176
  route_after_decision,
@@ -180,13 +217,16 @@ def build_workflow(
180
  },
181
  )
182
 
 
 
 
 
183
  graph.add_edge(
184
  "link",
185
  "complete",
186
  )
187
 
188
- # Temporary review branch returns to completion.
189
- # We will replace this with the real human gate later.
190
  graph.add_edge(
191
  "human_review",
192
  "complete",
@@ -197,7 +237,11 @@ def build_workflow(
197
  END,
198
  )
199
 
 
 
 
 
200
  return graph.compile(
201
  checkpointer=checkpointer,
202
  interrupt_before=interrupt_before,
203
- )
 
2
 
3
  from sqlalchemy.orm import Session
4
  from langgraph.graph import END, START, StateGraph
5
+ from langgraph.types import interrupt
6
 
7
  from app.llm.client import LLMClient
8
  from app.repositories.knowledge_repository import KnowledgeRepository
 
32
  ):
33
  """
34
  Build the DocWeave document intelligence workflow.
35
+
36
+ The human review branch uses LangGraph's durable interrupt()
37
+ mechanism so that the workflow state is checkpointed and can
38
+ later be resumed with Command(resume=...).
39
  """
40
 
41
  llm_client = LLMClient()
 
115
  ),
116
  )
117
 
118
+ # ------------------------------------------------------------------
119
+ # Durable human review gate
120
+ # ------------------------------------------------------------------
121
+
122
+ def human_review(
123
+ state: WorkflowState,
124
+ ) -> WorkflowState:
125
+ """
126
+ Durable human approval gate.
127
+
128
+ LangGraph checkpoints the current state when interrupt()
129
+ is reached. The workflow can later continue from this
130
+ checkpoint using Command(resume=...).
131
+ """
132
+
133
  state.current_node = "HUMAN_REVIEW"
134
+
135
+ interrupt(
136
+ {
137
+ "type": "human_review",
138
+ "workflow_run_id": str(
139
+ state.workflow_run_id
140
+ ),
141
+ "document_version_id": str(
142
+ state.document_version_id
143
+ ),
144
+ "message": (
145
+ "Human review is required before "
146
+ "the workflow can continue."
147
+ ),
148
+ }
149
+ )
150
+
151
  return state
152
 
153
  graph.add_node(
 
204
  "decision",
205
  )
206
 
207
+ # ------------------------------------------------------------------
208
+ # Decision routing
209
+ # ------------------------------------------------------------------
210
+
211
  graph.add_conditional_edges(
212
  "decision",
213
  route_after_decision,
 
217
  },
218
  )
219
 
220
+ # ------------------------------------------------------------------
221
+ # Normal completion path
222
+ # ------------------------------------------------------------------
223
+
224
  graph.add_edge(
225
  "link",
226
  "complete",
227
  )
228
 
229
+ # After human approval/resume, continue to completion.
 
230
  graph.add_edge(
231
  "human_review",
232
  "complete",
 
237
  END,
238
  )
239
 
240
+ # ------------------------------------------------------------------
241
+ # Compile
242
+ # ------------------------------------------------------------------
243
+
244
  return graph.compile(
245
  checkpointer=checkpointer,
246
  interrupt_before=interrupt_before,
247
+ )