From fbe8a7c38a57ff4d50b189eb031ea189c983f21d Mon Sep 17 00:00:00 2001 From: cristhianzl Date: Mon, 15 Jun 2026 18:10:11 -0300 Subject: [PATCH] human-in-the-loop pause/resume for v2 background runs --- src/backend/base/langflow/api/build.py | 73 ++++- .../base/langflow/api/v2/adapters/langflow.py | 1 + .../base/langflow/api/v2/agui_translator.py | 4 + src/backend/base/langflow/api/v2/hitl.py | 57 ++++ src/backend/base/langflow/api/v2/workflow.py | 66 +++- .../base/langflow/schema/content_block.py | 13 +- .../base/langflow/schema/content_types.py | 18 ++ .../services/background_execution/service.py | 17 +- .../base/langflow/services/jobs/service.py | 3 +- .../api/v2/adapters/test_human_input_event.py | 66 ++++ .../tests/unit/api/v2/test_resume_route.py | 95 ++++++ .../core/chatComponents/ContentDisplay.tsx | 23 +- .../core/chatComponents/HumanInputCard.tsx | 161 ++++++++++ .../__tests__/HumanInputCard.test.tsx | 155 +++++++++ .../chat-messages/components/bot-message.tsx | 26 +- .../consume-background-events.test.ts | 169 ++++++++++ .../agui/__tests__/run-flow-bridge.test.ts | 22 ++ .../src/controllers/API/agui/run-agent.ts | 33 ++ .../controllers/API/agui/run-flow-bridge.ts | 302 ++++++++++++++++++ .../queries/workflows/use-resume-workflow.ts | 43 +++ src/frontend/src/stores/flowStore.ts | 21 +- src/frontend/src/types/chat/index.ts | 29 +- src/frontend/src/types/zustand/flow/index.ts | 2 + src/lfx/src/lfx/graph/checkpoint/schema.py | 7 +- src/lfx/src/lfx/schema/content_block.py | 13 +- src/lfx/src/lfx/schema/content_types.py | 18 ++ src/lfx/src/lfx/schema/workflow.py | 15 + .../unit/graph/checkpoint/_static_pauser.py | 42 +++ .../checkpoint/test_resume_round_trip.py | 89 ++++++ .../unit/graph/checkpoint/test_schema.py | 15 + 30 files changed, 1563 insertions(+), 35 deletions(-) create mode 100644 src/backend/base/langflow/api/v2/hitl.py create mode 100644 src/backend/tests/unit/api/v2/adapters/test_human_input_event.py create mode 100644 src/backend/tests/unit/api/v2/test_resume_route.py create mode 100644 src/frontend/src/components/core/chatComponents/HumanInputCard.tsx create mode 100644 src/frontend/src/components/core/chatComponents/__tests__/HumanInputCard.test.tsx create mode 100644 src/frontend/src/controllers/API/agui/__tests__/consume-background-events.test.ts create mode 100644 src/frontend/src/controllers/API/queries/workflows/use-resume-workflow.ts create mode 100644 src/lfx/tests/unit/graph/checkpoint/_static_pauser.py create mode 100644 src/lfx/tests/unit/graph/checkpoint/test_resume_round_trip.py diff --git a/src/backend/base/langflow/api/build.py b/src/backend/base/langflow/api/build.py index 37b6ad6adf..920577e810 100644 --- a/src/backend/base/langflow/api/build.py +++ b/src/backend/base/langflow/api/build.py @@ -350,6 +350,43 @@ async def create_flow_response( ) +async def _persist_human_input_card(data: dict, flow_id: uuid.UUID, session_id: str, job_id) -> None: + """Persist the pause as a chat message so the interactive card survives reload. + + The card carries request_id + job_id, so a reloaded session can resume the run. + """ + from lfx.memory import astore_message + from lfx.schema.content_block import ContentBlock + from lfx.schema.content_types import HumanInputContent + from lfx.schema.message import Message + + content = HumanInputContent( + request_id=data.get("request_id", ""), + job_id=str(job_id) if job_id else None, + kind=data.get("kind", "node_input"), + prompt=data.get("prompt"), + options=data.get("options") or [], + fields=data.get("schema") or [], + allowed_decisions=data.get("allowed_decisions") or [], + ) + block = ContentBlock(title="Human input required", contents=[content]) + message = Message( + text="", + sender="Machine", + sender_name="AI", + session_id=session_id, + flow_id=flow_id, + content_blocks=[block], + ) + try: + stored = await astore_message(message, flow_id=flow_id, run_id=str(job_id) if job_id else None) + # Record the card's message id so resume can mark it answered (see the resume route). + if stored and job_id is not None: + await get_job_service().update_job_metadata(uuid.UUID(str(job_id)), {"card_message_id": str(stored[0].id)}) + except Exception: # noqa: BLE001 + await logger.awarning("Failed to persist human-input card for flow %s", flow_id, exc_info=True) + + async def generate_flow_events( *, flow_id: uuid.UUID, @@ -420,6 +457,12 @@ async def generate_flow_events( graph = await create_graph(fresh_session, flow_id_str, flow_name) graph.set_run_id(run_id) + if job_id is not None: + from lfx.services.deps import get_checkpoint_service + + graph.job_id = str(job_id) # the checkpoint is keyed by job_id + graph.checkpointing_enabled = True # arm the pause seam for producers + graph.checkpoint_store = get_checkpoint_service() first_layer = sort_vertices(graph) for vertex_id in first_layer: @@ -718,12 +761,14 @@ async def generate_flow_events( _build_run_id = uuid.UUID(graph.run_id) if graph.run_id else None if _build_run_id is not None: _build_job_svc = get_job_service() - await _build_job_svc.create_job( - job_id=_build_run_id, - flow_id=flow_id, - user_id=current_user.id, - job_type=JobType.WORKFLOW, - ) + # Background path already created the job; re-creating it = UNIQUE violation. + if await _build_job_svc.get_job_by_job_id(_build_run_id) is None: + await _build_job_svc.create_job( + job_id=_build_run_id, + flow_id=flow_id, + user_id=current_user.id, + job_type=JobType.WORKFLOW, + ) except Exception: # noqa: BLE001 await logger.awarning( "Failed to create workflow job for /build — memory base tracking disabled for flow %s", @@ -774,10 +819,18 @@ async def generate_flow_events( event_manager.on_error(data=error_message.data) raise - if _build_job_svc and _build_run_id: - await _build_job_svc.execute_with_status(_build_run_id, _run_vertex_build) - else: - await _run_vertex_build() + try: + runner_owns_status = job_id is not None # background path: JobRunner already wraps execute_with_status + if _build_job_svc and _build_run_id and not runner_owns_status: + await _build_job_svc.execute_with_status(_build_run_id, _run_vertex_build) + else: + await _run_vertex_build() + except GraphPausedException as exc: + # Non-terminal: persist the card to history, emit the pause event, end without on_end. + await _persist_human_input_card(exc.data or {}, flow_id, graph.session_id or str(flow_id), job_id) + event_manager.send_event(event_type="human_input_required", data=exc.data or {}) + await event_manager.queue.put((None, None, time.time())) + return build_duration = sum(vertex_timedeltas) event_manager.on_end(data={"build_duration": build_duration}) diff --git a/src/backend/base/langflow/api/v2/adapters/langflow.py b/src/backend/base/langflow/api/v2/adapters/langflow.py index 4a0417eb04..0df03340c9 100644 --- a/src/backend/base/langflow/api/v2/adapters/langflow.py +++ b/src/backend/base/langflow/api/v2/adapters/langflow.py @@ -34,6 +34,7 @@ _LANGFLOW_DURABLE_EVENTS: frozenset[str] = frozenset( "remove_message", "error", "end", + "human_input_required", } ) diff --git a/src/backend/base/langflow/api/v2/agui_translator.py b/src/backend/base/langflow/api/v2/agui_translator.py index b8db8e7120..f39f8597bd 100644 --- a/src/backend/base/langflow/api/v2/agui_translator.py +++ b/src/backend/base/langflow/api/v2/agui_translator.py @@ -124,6 +124,10 @@ class AGUITranslator: del self._buffered_messages[removed_id] self._emitted_text_message_ids.add(removed_id) return [CustomEvent(name="langflow.message.removed", value={"message_id": removed_id})] + if event_type == "human_input_required": + # Non-terminal: the run suspends for human input. Must precede the end/error + # branches so it never closes the open message or emits RUN_FINISHED/RUN_ERROR. + return [CustomEvent(name="langflow.human_input_required", value=data)] # Only terminal events close an open text message. Non-terminal events # (build_start, end_vertex, log, ...) interleave with tokens of the same diff --git a/src/backend/base/langflow/api/v2/hitl.py b/src/backend/base/langflow/api/v2/hitl.py new file mode 100644 index 0000000000..dcf350e72d --- /dev/null +++ b/src/backend/base/langflow/api/v2/hitl.py @@ -0,0 +1,57 @@ +"""Human-in-the-loop persistence helpers for the v2 workflows API. + +The pause is stored as a chat message so the interactive card survives reload; +on resume the same message is updated with the chosen action so a reloaded +session renders it as resolved instead of re-offering the decision. +""" + +from __future__ import annotations + +from uuid import UUID + +from sqlalchemy.orm.attributes import flag_modified + +from langflow.services.deps import get_job_service +from lfx.log import logger + + +def _set_submitted_action(content_blocks: list, action_id: str | None) -> bool: + """Stamp ``submitted_action`` on the card's content in place, preserving the rest.""" + changed = False + for block in content_blocks or []: + contents = block.get("contents") if isinstance(block, dict) else getattr(block, "contents", None) + for content in contents or []: + ctype = content.get("type") if isinstance(content, dict) else getattr(content, "type", None) + if ctype != "human_input": + continue + if isinstance(content, dict): + content["submitted_action"] = action_id + else: + content.submitted_action = action_id + changed = True + return changed + + +async def mark_card_answered(job_id: UUID, request_id: str, decision: dict) -> None: # noqa: ARG001 + """Record the chosen action on the persisted card message for ``job_id``. + + Patches the existing card in place so the prompt/options it was stored with are + preserved — rebuilding from the suspend event would drop them. + """ + from langflow.services.database.models.message.model import MessageTable + from lfx.services.deps import session_scope + + job = await get_job_service().get_job_by_job_id(job_id) + card_message_id = (job.job_metadata or {}).get("card_message_id") if job else None + if not card_message_id: + return + try: + async with session_scope() as session: + message = await session.get(MessageTable, UUID(str(card_message_id))) + if message is None: + return + if _set_submitted_action(message.content_blocks, decision.get("action_id")): + flag_modified(message, "content_blocks") + session.add(message) + except Exception as _e: # noqa: BLE001 + await logger.awarning("Failed to mark human-input card answered for job %s: %r", job_id, _e) diff --git a/src/backend/base/langflow/api/v2/workflow.py b/src/backend/base/langflow/api/v2/workflow.py index 0416d6533d..4e718e822c 100644 --- a/src/backend/base/langflow/api/v2/workflow.py +++ b/src/backend/base/langflow/api/v2/workflow.py @@ -44,6 +44,8 @@ from lfx.schema.workflow import ( JobStatus, WorkflowExecutionResponse, WorkflowJobResponse, + WorkflowResumeRequest, + WorkflowResumeResponse, WorkflowRunRequest, WorkflowStopRequest, WorkflowStopResponse, @@ -688,7 +690,11 @@ async def _stream_event_frames( ) seq += 1 for event in adapter.translate(event_type, event_data): - yield _frame(event, seq) + frame_bytes, frame_type = _frame(event, seq) + # Runner detects a pause by the langflow-side type; agui maps it to CUSTOM. + if event_type == "human_input_required": + frame_type = "human_input_required" + yield (frame_bytes, frame_type) seq += 1 for event in adapter.final_events(): yield _frame(event, seq) @@ -1152,6 +1158,64 @@ async def stop_workflow( ) from exc +@router.post( + "/{job_id}/resume", + summary="Resume Workflow", + description="Resume a suspended (human-in-the-loop) workflow with a decision.", +) +async def resume_workflow( + job_id: str, + request: WorkflowResumeRequest, + current_user: Annotated[UserRead, Depends(get_current_user_for_workflow)], +) -> WorkflowResumeResponse: + """Resume a SUSPENDED workflow run with a human decision. + + Owner-or-superuser; a non-owner non-superuser (or unknown/non-workflow job) + maps to 404 to avoid leaking other users' runs. A stale/duplicate request_id + or a non-suspended job maps to 409 (single-use enforced behind ``resume_job``). + """ + + def _not_found() -> HTTPException: + return HTTPException( + status_code=status.HTTP_404_NOT_FOUND, + detail={"error": "Workflow job not found", "code": "JOB_NOT_FOUND", "job_id": job_id}, + ) + + try: + parsed_job_id = UUID(job_id) + except ValueError as exc: + raise _not_found() from exc + + job = await get_job_service().get_job_by_job_id(parsed_job_id) + is_owner = job is not None and (job.user_id is None or job.user_id == current_user.id) + if job is None or job.type != JobType.WORKFLOW or not (is_owner or current_user.is_superuser): + raise _not_found() + + service = get_background_execution_service() + if service._frame_source_factory is None: # noqa: SLF001 + service._frame_source_factory = _default_frame_source_factory # noqa: SLF001 + accepted = await service.resume_job( + parsed_job_id, + current_user, + request_id=request.request_id, + decision=request.decision or {}, + ) + if not accepted: + raise HTTPException( + status_code=status.HTTP_409_CONFLICT, + detail={ + "error": "Job is not resumable", + "code": "NOT_RESUMABLE", + "message": "Job is not suspended, already resumed, or the request_id is stale.", + "job_id": job_id, + }, + ) + from langflow.api.v2.hitl import mark_card_answered + + await mark_card_answered(parsed_job_id, request.request_id, request.decision or {}) + return WorkflowResumeResponse(job_id=job_id, status="resuming", message="Resume accepted") + + @router.get( "/{job_id}/events", response_model=None, diff --git a/src/backend/base/langflow/schema/content_block.py b/src/backend/base/langflow/schema/content_block.py index 6df7ef6e8b..bbe266cc41 100644 --- a/src/backend/base/langflow/schema/content_block.py +++ b/src/backend/base/langflow/schema/content_block.py @@ -3,7 +3,15 @@ from typing import Annotated from pydantic import BaseModel, Discriminator, Field, Tag, field_serializer, field_validator from typing_extensions import TypedDict -from .content_types import CodeContent, ErrorContent, JSONContent, MediaContent, TextContent, ToolContent +from .content_types import ( + CodeContent, + ErrorContent, + HumanInputContent, + JSONContent, + MediaContent, + TextContent, + ToolContent, +) def _get_type(d: dict | BaseModel) -> str | None: @@ -19,7 +27,8 @@ ContentType = Annotated[ | Annotated[TextContent, Tag("text")] | Annotated[MediaContent, Tag("media")] | Annotated[CodeContent, Tag("code")] - | Annotated[JSONContent, Tag("json")], + | Annotated[JSONContent, Tag("json")] + | Annotated[HumanInputContent, Tag("human_input")], Discriminator(_get_type), ] diff --git a/src/backend/base/langflow/schema/content_types.py b/src/backend/base/langflow/schema/content_types.py index d3f580a36d..449b1471ae 100644 --- a/src/backend/base/langflow/schema/content_types.py +++ b/src/backend/base/langflow/schema/content_types.py @@ -78,6 +78,24 @@ class CodeContent(BaseContent): title: str | None = None +class HumanInputContent(BaseContent): + """Content type for a human-in-the-loop pause persisted in the chat history. + + Carries the pending decision so the interactive card survives reload: the + request_id and job_id let the card resume the suspended run after an F5. + """ + + type: Literal["human_input"] = Field(default="human_input") + request_id: str + job_id: str | None = None + kind: str = "node_input" + prompt: str | None = None + options: list[dict[str, Any]] = Field(default_factory=list) + fields: list[dict[str, Any]] = Field(default_factory=list) + allowed_decisions: list[str] = Field(default_factory=list) + submitted_action: str | None = None + + class ToolContent(BaseContent): """Content type for tool start content.""" diff --git a/src/backend/base/langflow/services/background_execution/service.py b/src/backend/base/langflow/services/background_execution/service.py index 1a36d20eac..877b13d18f 100644 --- a/src/backend/base/langflow/services/background_execution/service.py +++ b/src/backend/base/langflow/services/background_execution/service.py @@ -225,8 +225,12 @@ class BackgroundExecutionService(Service): return async def _is_terminal() -> bool: + # SUSPENDED ends the tail too: a run that connected while IN_PROGRESS and + # then suspended has no live tail to wait on (the pause isn't published live). current = await job_service.get_job_by_job_id(job_id) - return current is not None and current.status in _TERMINAL_STATUSES + return current is not None and ( + current.status in _TERMINAL_STATUSES or current.status == JobStatus.SUSPENDED + ) async for frame in self._bus.reattach( str(job_id), last_seq=last_seq, read_durable=read_durable, is_done=_is_terminal @@ -266,20 +270,21 @@ class BackgroundExecutionService(Service): await job_service.write_signal(job_id, SignalType.STOP) await self._executor.cancel(str(job_id)) - async def resume_job(self, job_id: UUID, user: UserRead, *, request_id: str, decision: Any) -> bool: + async def resume_job(self, job_id: UUID, user: UserRead, *, request_id: str, decision: Any) -> bool: # noqa: ARG002 """Carry a human decision back into a SUSPENDED run and re-enqueue it. Returns True when the run was accepted for resume, False on a conflict (not suspended, stale request_id, or lost the single-flight flip) — the - route maps False to 409. Ownership is enforced by ``_validate``. + route maps False to 409. Ownership (owner-or-superuser) is enforced at the + HTTP route, so this fetches by id alone and trusts the already-validated user. """ - job = await self._validate(job_id, user) - if job.status != JobStatus.SUSPENDED: + job_service = get_job_service() + job = await job_service.get_job_by_job_id(job_id) + if job is None or job.type != JobType.WORKFLOW or job.status != JobStatus.SUSPENDED: return False pending = (job.job_metadata or {}).get("pending_request_id") if pending is not None and request_id != pending: return False - job_service = get_job_service() # Win the single-flight flip BEFORE writing the RESUME signal, so exactly one # RESUME row exists per suspend and a loser never strands a stray decision. if not await job_service.claim_suspended_for_resume(job_id, owner=self._owner): diff --git a/src/backend/base/langflow/services/jobs/service.py b/src/backend/base/langflow/services/jobs/service.py index 121f51ea3e..41d228a387 100644 --- a/src/backend/base/langflow/services/jobs/service.py +++ b/src/backend/base/langflow/services/jobs/service.py @@ -10,6 +10,7 @@ if TYPE_CHECKING: from datetime import datetime, timezone from uuid import UUID +from lfx.graph.exceptions import GraphPausedException from sqlalchemy.exc import IntegrityError, OperationalError from sqlmodel import col, func, select @@ -751,7 +752,7 @@ class JobService(Service): await logger.ainfo(f"Executing job function for job_id={job_id}") result = await run_coro_func(*args, **kwargs) - except PauseRequested: + except (PauseRequested, GraphPausedException): # A producer paused the run for human input. The runner suspends it # (SUSPENDED); do NOT write a terminal status here. await logger.adebug(f"Job {job_id} paused for human input") diff --git a/src/backend/tests/unit/api/v2/adapters/test_human_input_event.py b/src/backend/tests/unit/api/v2/adapters/test_human_input_event.py new file mode 100644 index 0000000000..f0a06f29e6 --- /dev/null +++ b/src/backend/tests/unit/api/v2/adapters/test_human_input_event.py @@ -0,0 +1,66 @@ +"""HITL human_input_required translation on both stream protocols (LE-1450). + +The pause is a non-terminal durable milestone: langflow passes it through and +marks it durable; agui rides the existing CUSTOM mechanism as a single +``langflow.human_input_required`` event that never closes the message or ends the run. +""" + +from __future__ import annotations + +import json + +from langflow.api.v2.adapters import StreamAdapterContext, get_stream_adapter + +_PAYLOAD = { + "request_id": "node:job-1", + "kind": "node_input", + "prompt": "Approve refund?", + "options": [{"action_id": "approve", "label": "Approve"}, {"action_id": "reject", "label": "Reject"}], + "schema": [], + "allowed_decisions": ["approve", "reject"], +} + + +def _ctx() -> StreamAdapterContext: + return StreamAdapterContext(run_id="run-1", thread_id="thread-1") + + +class TestLangflowHumanInputEvent: + def test_translate_passes_through_one_frame(self): + adapter = get_stream_adapter("langflow", _ctx()) + events = list(adapter.translate("human_input_required", _PAYLOAD)) + assert len(events) == 1 + payload = json.loads(events[0].data_json) + assert payload == {"event": "human_input_required", "data": _PAYLOAD} + + def test_is_durable_true(self): + adapter = get_stream_adapter("langflow", _ctx()) + assert adapter.is_durable("human_input_required") is True + + def test_not_the_terminal_error_type(self): + adapter = get_stream_adapter("langflow", _ctx()) + assert adapter.terminal_error_type != "human_input_required" + + +class TestAguiHumanInputEvent: + def test_translate_emits_single_custom_event(self): + adapter = get_stream_adapter("agui", _ctx()) + events = list(adapter.translate("human_input_required", _PAYLOAD)) + assert len(events) == 1 + assert events[0].type == "CUSTOM" + payload = json.loads(events[0].data_json) + assert payload["name"] == "langflow.human_input_required" + assert payload["value"] == _PAYLOAD + + def test_is_durable_via_custom(self): + adapter = get_stream_adapter("agui", _ctx()) + assert adapter.is_durable("CUSTOM") is True + + def test_pause_does_not_end_the_run_or_close_the_message(self): + adapter = get_stream_adapter("agui", _ctx()) + list(adapter.translate("token", {"id": "m1", "chunk": "Working"})) # open a text message + events = list(adapter.translate("human_input_required", _PAYLOAD)) + types = [e.type for e in events] + assert "RUN_FINISHED" not in types + assert "RUN_ERROR" not in types + assert "TEXT_MESSAGE_END" not in types # the open message stays open across the pause diff --git a/src/backend/tests/unit/api/v2/test_resume_route.py b/src/backend/tests/unit/api/v2/test_resume_route.py new file mode 100644 index 0000000000..8fc6dc0ddc --- /dev/null +++ b/src/backend/tests/unit/api/v2/test_resume_route.py @@ -0,0 +1,95 @@ +"""HTTP resume route (LE-1450): POST /api/v2/workflows/{job_id}/resume. + +Route-level tests against a hand-written SUSPENDED job row (the full HTTP +round-trip gates on a real pausing flow): owner-or-superuser auth with +deny-to-404, single-use 409, and accepted-resume 200. +""" + +from __future__ import annotations + +from uuid import uuid4 + +import pytest +from langflow.services.database.models.flow.model import Flow +from langflow.services.database.models.jobs.model import Job, JobStatus, JobType +from lfx.services.deps import session_scope + +pytestmark = pytest.mark.usefixtures("client") + + +def _headers(api_key) -> dict: + return {"x-api-key": api_key.api_key} + + +@pytest.fixture +async def suspended_job(created_api_key): + """A hand-written SUSPENDED workflow job owned by the api-key user, with a real flow.""" + user_id = created_api_key.user_id + flow_id, job_id = uuid4(), uuid4() + request = {"flow_id": str(flow_id), "mode": "background", "stream_protocol": "langflow", "input_value": "hi"} + async with session_scope() as session: + session.add(Flow(id=flow_id, name=f"f-{flow_id}", data={"nodes": [], "edges": []}, user_id=user_id)) + session.add( + Job( + job_id=job_id, + flow_id=flow_id, + user_id=user_id, + type=JobType.WORKFLOW, + status=JobStatus.SUSPENDED, + job_metadata={"pending_request_id": "req-1", "request": request}, + ) + ) + await session.flush() + yield job_id + async with session_scope() as session: + job = await session.get(Job, job_id) + if job: + await session.delete(job) + flow = await session.get(Flow, flow_id) + if flow: + await session.delete(flow) + + +async def test_resume_owner_accepts_200(client, created_api_key, suspended_job): + body = {"request_id": "req-1", "decision": {"action_id": "approve"}} + resp = await client.post(f"api/v2/workflows/{suspended_job}/resume", json=body, headers=_headers(created_api_key)) + assert resp.status_code == 200 + assert resp.json()["status"] == "resuming" + + +async def test_resume_stale_request_id_409(client, created_api_key, suspended_job): + body = {"request_id": "WRONG", "decision": {"action_id": "approve"}} + resp = await client.post(f"api/v2/workflows/{suspended_job}/resume", json=body, headers=_headers(created_api_key)) + assert resp.status_code == 409 + + +async def test_resume_unknown_job_404(client, created_api_key): + body = {"request_id": "req-1", "decision": {}} + resp = await client.post(f"api/v2/workflows/{uuid4()}/resume", json=body, headers=_headers(created_api_key)) + assert resp.status_code == 404 + + +async def test_resume_non_suspended_job_409(client, created_api_key): + """A QUEUED (not suspended) job is not resumable → 409.""" + user_id = created_api_key.user_id + flow_id, job_id = uuid4(), uuid4() + async with session_scope() as session: + session.add(Flow(id=flow_id, name=f"f-{flow_id}", data={"nodes": [], "edges": []}, user_id=user_id)) + session.add(Job(job_id=job_id, flow_id=flow_id, user_id=user_id, type=JobType.WORKFLOW, status=JobStatus.QUEUED)) + await session.flush() + try: + body = {"request_id": "req-1", "decision": {}} + resp = await client.post(f"api/v2/workflows/{job_id}/resume", json=body, headers=_headers(created_api_key)) + assert resp.status_code == 409 + finally: + async with session_scope() as session: + for model, key in ((Job, job_id), (Flow, flow_id)): + row = await session.get(model, key) + if row: + await session.delete(row) + + +async def test_resume_invalid_job_id_404(client, created_api_key): + body = {"request_id": "req-1", "decision": {}} + resp = await client.post("api/v2/workflows/not-a-uuid/resume", json=body, headers=_headers(created_api_key)) + assert resp.status_code == 404 diff --git a/src/frontend/src/components/core/chatComponents/ContentDisplay.tsx b/src/frontend/src/components/core/chatComponents/ContentDisplay.tsx index 23d3ffcc92..2933aeb1af 100644 --- a/src/frontend/src/components/core/chatComponents/ContentDisplay.tsx +++ b/src/frontend/src/components/core/chatComponents/ContentDisplay.tsx @@ -2,20 +2,28 @@ import type { ReactNode } from "react"; import Markdown from "react-markdown"; import rehypeMathjax from "rehype-mathjax/browser"; import remarkGfm from "remark-gfm"; -import type { ContentType, JSONValue } from "@/types/chat"; +import type { ContentType, InteractiveContent, JSONValue } from "@/types/chat"; import { extractLanguage, isCodeBlock } from "@/utils/codeBlockUtils"; import ForwardedIconComponent from "../../common/genericIconComponent"; import SimplifiedCodeTabComponent from "../codeTabsComponent"; import DurationDisplay from "./DurationDisplay"; +import HumanInputCard, { type HumanInputDecision } from "./HumanInputCard"; export default function ContentDisplay({ content, chatId, playgroundPage, + humanInputSubmitted, + onHumanInputSubmit, }: { content: ContentType; chatId: string; playgroundPage?: boolean; + humanInputSubmitted?: boolean; + onHumanInputSubmit?: ( + content: InteractiveContent, + decision: HumanInputDecision, + ) => void; }) { const renderDuration = content.duration !== undefined && !playgroundPage && (
@@ -242,6 +250,19 @@ export default function ContentDisplay({
); break; + case "human_input": + contentData = ( + onHumanInputSubmit(content, decision) + : undefined + } + /> + ); + break; } return ( diff --git a/src/frontend/src/components/core/chatComponents/HumanInputCard.tsx b/src/frontend/src/components/core/chatComponents/HumanInputCard.tsx new file mode 100644 index 0000000000..ce11072cce --- /dev/null +++ b/src/frontend/src/components/core/chatComponents/HumanInputCard.tsx @@ -0,0 +1,161 @@ +import { useState } from "react"; +import Markdown from "react-markdown"; +import remarkGfm from "remark-gfm"; +import { queryClient } from "@/contexts"; +import { + consumeBackgroundEvents, + getResumeContext, + markHumanInputSubmitted, +} from "@/controllers/API/agui/run-flow-bridge"; +import { useResumeWorkflow } from "@/controllers/API/queries/workflows/use-resume-workflow"; +import { Button } from "@/components/ui/button"; +import { Input } from "@/components/ui/input"; +import useAlertStore from "@/stores/alertStore"; +import useFlowStore from "@/stores/flowStore"; +import type { InteractiveContent } from "@/types/chat"; +import ForwardedIconComponent from "../../common/genericIconComponent"; + +export interface HumanInputDecision { + action_id: string; + values: Record; +} + +/** + * Interactive human-in-the-loop card: a prompt, optional editable form fields, and + * one button per decision. Mirrors the convergent HITL pattern (a paused turn the + * user resolves inline) using Langflow's own components. Submit is single-use in the + * UI (best-effort; the server is authoritative and returns 409 on a stale resume). + */ +export default function HumanInputCard({ + content, + onSubmit, + submitted = false, +}: { + content: InteractiveContent; + onSubmit?: (decision: HumanInputDecision) => void; + submitted?: boolean; +}) { + const [values, setValues] = useState>({}); + const [localChosen, setLocalChosen] = useState(null); + const fields = content.schema ?? content.fields ?? []; + const { mutate: resume, isPending } = useResumeWorkflow(); + const setErrorData = useAlertStore((s) => s.setErrorData); + // Derive from the persisted/cached content first so the choice survives re-renders + // (the resume reattach replays the stream and re-renders the message list). + const chosen = content.submitted_action ?? localChosen; + const isSubmitted = submitted || chosen !== null; + + // Resume the suspended run and reattach to its durable event stream so the + // continued run renders into the same chat session. + const resumeRun = (decision: HumanInputDecision) => { + const jobId = content.job_id; + if (!jobId) { + setErrorData({ title: "Cannot resume: missing job id" }); + setLocalChosen(null); + return; + } + resume( + { jobId, requestId: content.request_id, decision }, + { + onSuccess: () => { + const flowStore = useFlowStore.getState(); + flowStore.setAwaitingInput(false); + const reattach = getResumeContext(content.request_id); + if (reattach) { + flowStore.setIsBuilding(true); + void consumeBackgroundEvents(reattach.jobId, reattach.opts); + } else { + // After a reload there is no live stream to reattach to; refetch the + // history so the resumed run's persisted output shows up. + void queryClient.invalidateQueries({ + queryKey: ["useGetMessagesQuery"], + }); + } + }, + onError: (err: Error) => { + // 409 means the run was already resumed (single-use) — keep the card + // locked on the choice; only a genuine failure re-opens the buttons. + const status = (err as { response?: { status?: number } })?.response + ?.status; + if (status === 409) return; + setLocalChosen(null); + setErrorData({ title: "Resume failed", list: [err.message] }); + }, + }, + ); + }; + + const handleDecision = (actionId: string) => { + if (isSubmitted || isPending) return; + setLocalChosen(actionId); + markHumanInputSubmitted(content.request_id, actionId); + const decision = { action_id: actionId, values }; + if (onSubmit) onSubmit(decision); + else resumeRun(decision); + }; + + return ( +
+
+ + Human input required +
+ + {content.prompt && ( +
+ {content.prompt} +
+ )} + + {fields.length > 0 && ( +
+ {fields.map((field) => ( + + ))} +
+ )} + +
+ {content.options + .filter((option) => chosen === null || option.action_id === chosen) + .map((option) => ( + + ))} +
+
+ ); +} diff --git a/src/frontend/src/components/core/chatComponents/__tests__/HumanInputCard.test.tsx b/src/frontend/src/components/core/chatComponents/__tests__/HumanInputCard.test.tsx new file mode 100644 index 0000000000..a58fb9a7b5 --- /dev/null +++ b/src/frontend/src/components/core/chatComponents/__tests__/HumanInputCard.test.tsx @@ -0,0 +1,155 @@ +import { fireEvent, render, screen } from "@testing-library/react"; +import type { InteractiveContent } from "@/types/chat"; +import HumanInputCard from "../HumanInputCard"; + +const mockResume = jest.fn(); +const mockConsume = jest.fn(); +const mockSetAwaitingInput = jest.fn(); +const mockSetIsBuilding = jest.fn(); +const mockSetErrorData = jest.fn(); + +jest.mock("@/contexts", () => ({ + queryClient: { invalidateQueries: jest.fn() }, +})); +jest.mock("@/controllers/API/queries/workflows/use-resume-workflow", () => ({ + useResumeWorkflow: () => ({ mutate: mockResume, isPending: false }), +})); +jest.mock("@/controllers/API/agui/run-flow-bridge", () => ({ + consumeBackgroundEvents: (...args: unknown[]) => mockConsume(...args), + getResumeContext: () => ({ jobId: "job-1", opts: { flowId: "f1" } }), + markHumanInputSubmitted: jest.fn(), +})); +jest.mock("@/stores/flowStore", () => ({ + __esModule: true, + default: { + getState: () => ({ + setAwaitingInput: mockSetAwaitingInput, + setIsBuilding: mockSetIsBuilding, + }), + }, +})); +jest.mock("@/stores/alertStore", () => ({ + __esModule: true, + default: (selector: (s: { setErrorData: unknown }) => unknown) => + selector({ setErrorData: mockSetErrorData }), +})); + +jest.mock("../../../common/genericIconComponent", () => ({ + __esModule: true, + default: ({ name }: { name: string }) => , +})); + +jest.mock("react-markdown", () => ({ + __esModule: true, + default: ({ children }: { children: string }) => {children}, +})); +jest.mock("remark-gfm", () => ({ __esModule: true, default: () => {} })); + +jest.mock("@/components/ui/button", () => ({ + Button: ({ children, ...props }: any) => , +})); +jest.mock("@/components/ui/input", () => ({ + Input: (props: any) => , +})); + +const _approval: InteractiveContent = { + type: "human_input", + kind: "tool_approval", + request_id: "node:job-1", + prompt: "Approve refund?", + options: [ + { action_id: "approve", label: "Approve" }, + { action_id: "reject", label: "Reject" }, + ], + allowed_decisions: ["approve", "reject"], +}; + +describe("HumanInputCard", () => { + beforeEach(() => jest.clearAllMocks()); + + it("renders the prompt and one button per option", () => { + render(); + expect(screen.getByText("Approve refund?")).toBeInTheDocument(); + expect(screen.getByTestId("human-input-decision-approve")).toBeInTheDocument(); + expect(screen.getByTestId("human-input-decision-reject")).toBeInTheDocument(); + }); + + it("submits the chosen action with empty values when there are no form fields", () => { + const onSubmit = jest.fn(); + render(); + fireEvent.click(screen.getByTestId("human-input-decision-approve")); + expect(onSubmit).toHaveBeenCalledWith({ action_id: "approve", values: {} }); + }); + + it("disables the controls once submitted", () => { + const onSubmit = jest.fn(); + render(); + fireEvent.click(screen.getByTestId("human-input-decision-approve")); + expect(onSubmit).not.toHaveBeenCalled(); + expect(screen.getByTestId("human-input-decision-approve")).toBeDisabled(); + }); + + it("collects form field values into the decision (node_input)", () => { + const onSubmit = jest.fn(); + const nodeInput: InteractiveContent = { + type: "human_input", + kind: "node_input", + request_id: "node:job-2", + prompt: "Pick one", + options: [ + { action_id: "a", label: "A" }, + { action_id: "b", label: "B" }, + { action_id: "c", label: "C" }, + ], + schema: [{ name: "reason", type: "str", required: true }], + allowed_decisions: ["a", "b", "c"], + }; + render(); + expect(screen.getAllByRole("button")).toHaveLength(3); + fireEvent.change(screen.getByTestId("human-input-field-reason"), { + target: { value: "fraud" }, + }); + fireEvent.click(screen.getByTestId("human-input-decision-b")); + expect(onSubmit).toHaveBeenCalledWith({ + action_id: "b", + values: { reason: "fraud" }, + }); + }); + + it("resumes the run itself when no onSubmit is provided", () => { + const content: InteractiveContent = { ..._approval, job_id: "job-1" }; + render(); + fireEvent.click(screen.getByTestId("human-input-decision-approve")); + expect(mockResume).toHaveBeenCalledWith( + { + jobId: "job-1", + requestId: "node:job-1", + decision: { action_id: "approve", values: {} }, + }, + expect.objectContaining({ onSuccess: expect.any(Function) }), + ); + const { onSuccess } = mockResume.mock.calls[0][1]; + onSuccess(); + expect(mockSetAwaitingInput).toHaveBeenCalledWith(false); + expect(mockSetIsBuilding).toHaveBeenCalledWith(true); + expect(mockConsume).toHaveBeenCalledWith("job-1", { flowId: "f1" }); + }); + + it("disables controls after a self-resume click (single-use)", () => { + const content: InteractiveContent = { ..._approval, job_id: "job-1" }; + render(); + fireEvent.click(screen.getByTestId("human-input-decision-approve")); + expect(screen.getByTestId("human-input-decision-approve")).toBeDisabled(); + }); + + it("keeps only the chosen option and removes the others after selecting", () => { + const content: InteractiveContent = { ..._approval, job_id: "job-1" }; + render(); + expect(screen.getByTestId("human-input-decision-reject")).toBeInTheDocument(); + fireEvent.click(screen.getByTestId("human-input-decision-approve")); + expect(screen.getByTestId("human-input-decision-approve")).toBeInTheDocument(); + expect( + screen.queryByTestId("human-input-decision-reject"), + ).not.toBeInTheDocument(); + }); +}); diff --git a/src/frontend/src/components/core/playgroundComponent/chat-view/chat-messages/components/bot-message.tsx b/src/frontend/src/components/core/playgroundComponent/chat-view/chat-messages/components/bot-message.tsx index 4f03f6b248..ca175b9e2b 100644 --- a/src/frontend/src/components/core/playgroundComponent/chat-view/chat-messages/components/bot-message.tsx +++ b/src/frontend/src/components/core/playgroundComponent/chat-view/chat-messages/components/bot-message.tsx @@ -6,7 +6,9 @@ import IconComponent, { } from "@/components/common/genericIconComponent"; import MessageMetadata from "@/components/common/messageMetadataComponent"; import { ContentBlockDisplay } from "@/components/core/chatComponents/ContentBlockDisplay"; +import HumanInputCard from "@/components/core/chatComponents/HumanInputCard"; import { useUpdateMessage } from "@/controllers/API/queries/messages"; +import type { InteractiveContent } from "@/types/chat"; import { CustomMarkdownField } from "@/customization/components/custom-markdown-field"; import useAlertStore from "@/stores/alertStore"; import useFlowStore from "@/stores/flowStore"; @@ -44,6 +46,12 @@ export const BotMessage = memo( const isEmpty = decodedMessage?.trim() === ""; const chatMessage = chat.message ? chat.message.toString() : ""; + // ContentBlockDisplay renders only tool_use content, so the HITL card is rendered directly below. + const humanInputContent = chat.content_blocks + ?.flatMap((block) => block.contents ?? []) + .find((content) => content?.type === "human_input") as + | InteractiveContent + | undefined; const { mutate: updateMessageMutation } = useUpdateMessage(); const handleEditMessage = (message: string) => { @@ -138,7 +146,10 @@ export const BotMessage = memo( )} >
- {(thinkingActive || displayTime > 0 || chatMessage !== "") && ( + {(thinkingActive || + displayTime > 0 || + chatMessage !== "" || + humanInputContent) && (
+ {humanInputContent && ( +
+ +
+ )} + {((chat.content_blocks && chat.content_blocks.length > 0) || (isBuilding && lastMessage)) && ( - {(chatMessage === "" || (isEmpty && !isStreaming)) && - isBuilding && - lastMessage ? ( + {humanInputContent ? null : (chatMessage === "" || + (isEmpty && !isStreaming)) && + isBuilding && + lastMessage ? ( ({ + __esModule: true, + default: { + getState: () => ({ + setIsBuilding, + setAwaitingInput, + setBuildInfo, + updateEdgesRunningByNodes, + revertBuiltStatusFromBuilding, + }), + }, +})); +jest.mock("@/stores/messagesStore", () => ({ + useMessagesStore: { + getState: () => ({ + addMessage, + messages: storeMessages, + }), + }, +})); +jest.mock("@/stores/alertStore", () => ({ + __esModule: true, + default: { getState: () => ({ setErrorData }) }, +})); +jest.mock( + "@/components/core/playgroundComponent/chat-view/utils/message-event-handler", + () => ({ handleMessageEvent: jest.fn() }), +); +const updateMessageMock = jest.fn(); +jest.mock( + "@/components/core/playgroundComponent/chat-view/utils/message-utils", + () => ({ updateMessage: (...args: unknown[]) => updateMessageMock(...args) }), +); + +import { consumeBackgroundEvents } from "@/controllers/API/agui/run-flow-bridge"; + +function sseStream(frames: string[]): ReadableStream { + const encoder = new TextEncoder(); + return new ReadableStream({ + start(controller) { + for (const frame of frames) { + controller.enqueue(encoder.encode(`data: ${frame}\n\n`)); + } + controller.close(); + }, + }); +} + +/** Emits real `data:\nid:` frames and closes WITHOUT a trailing blank line. */ +function sseStreamNoTrailingBlank(frames: string[]): ReadableStream { + const encoder = new TextEncoder(); + const body = + frames.map((f, i) => `data: ${f}\nid: ${i + 1}`).join("\n\n") + "\n"; + return new ReadableStream({ + start(controller) { + controller.enqueue(encoder.encode(body)); + controller.close(); + }, + }); +} + +describe("consumeBackgroundEvents", () => { + beforeEach(() => { + jest.clearAllMocks(); + storeMessages = []; + }); + + it("injects the card and parks awaiting-input when the run suspends", async () => { + const humanInput = JSON.stringify({ + type: "CUSTOM", + name: "langflow.human_input_required", + value: { + request_id: "HI:job-1", + kind: "node_input", + prompt: "Approve?", + options: [{ action_id: "approve", label: "Approve" }], + schema: [], + allowed_decisions: ["approve"], + }, + }); + global.fetch = jest.fn().mockResolvedValue({ + ok: true, + body: sseStream([ + JSON.stringify({ type: "RUN_STARTED", runId: "job-1" }), + humanInput, + ]), + }) as unknown as typeof fetch; + + await consumeBackgroundEvents("job-1", { flowId: "f1", threadId: "s1" }); + + expect(addMessage).toHaveBeenCalledTimes(1); + expect(addMessage.mock.calls[0][0].id).toBe("human-input-HI:job-1"); + // The new playground reads the react-query cache, not useMessagesStore. + expect(updateMessageMock).toHaveBeenCalledTimes(1); + expect(updateMessageMock.mock.calls[0][0].id).toBe("human-input-HI:job-1"); + expect(setIsBuilding).toHaveBeenLastCalledWith(false); + expect(setAwaitingInput).toHaveBeenLastCalledWith(true); + }); + + it("processes the final pause frame even when the stream closes without a trailing blank line", async () => { + const humanInput = JSON.stringify({ + type: "CUSTOM", + name: "langflow.human_input_required", + value: { + request_id: "HI:job-2", + kind: "node_input", + prompt: "Approve?", + options: [{ action_id: "approve", label: "Approve" }], + schema: [], + allowed_decisions: ["approve"], + }, + }); + global.fetch = jest.fn().mockResolvedValue({ + ok: true, + body: sseStreamNoTrailingBlank([ + JSON.stringify({ type: "RUN_STARTED", runId: "job-2" }), + JSON.stringify({ type: "STATE_DELTA", delta: [] }), + humanInput, + ]), + }) as unknown as typeof fetch; + + await consumeBackgroundEvents("job-2", { flowId: "f1", threadId: "s1" }); + + expect(addMessage).toHaveBeenCalledTimes(1); + expect(addMessage.mock.calls[0][0].id).toBe("human-input-HI:job-2"); + expect(setAwaitingInput).toHaveBeenLastCalledWith(true); + }); + + it("clears awaiting-input when a resumed run reaches RUN_FINISHED", async () => { + global.fetch = jest.fn().mockResolvedValue({ + ok: true, + body: sseStream([ + JSON.stringify({ type: "RUN_STARTED", runId: "job-1" }), + JSON.stringify({ type: "RUN_FINISHED" }), + ]), + }) as unknown as typeof fetch; + + await consumeBackgroundEvents("job-1", { flowId: "f1", threadId: "s1" }); + + expect(setIsBuilding).toHaveBeenLastCalledWith(false); + expect(setAwaitingInput).toHaveBeenLastCalledWith(false); + expect(setBuildInfo).toHaveBeenCalledWith({ success: true }); + }); +}); diff --git a/src/frontend/src/controllers/API/agui/__tests__/run-flow-bridge.test.ts b/src/frontend/src/controllers/API/agui/__tests__/run-flow-bridge.test.ts index f61cd9f432..9c8ccf4ab2 100644 --- a/src/frontend/src/controllers/API/agui/__tests__/run-flow-bridge.test.ts +++ b/src/frontend/src/controllers/API/agui/__tests__/run-flow-bridge.test.ts @@ -20,6 +20,10 @@ function makeRecordingContext() { setRunId: (runId) => calls.push(`setRunId:${runId}`), applyDelta: (ops) => calls.push(`applyDelta:${ops.length}`), handleCustomEvent: (eventType) => calls.push(`custom:${eventType}`), + onHumanInput: (payload) => + calls.push( + `humanInput:${(payload as { request_id?: string })?.request_id ?? ""}`, + ), onFinished: () => calls.push("finished"), onError: (message) => calls.push(`error:${message}`), }; @@ -139,3 +143,21 @@ describe("handleAGUIEvent non-terminal contract", () => { expect(calls).toEqual([]); }); }); + +describe("handleAGUIEvent human-input contract", () => { + it("surfaces a human_input_required CUSTOM event non-terminally", () => { + const { ctx, calls } = makeRecordingContext(); + + const terminal = handleAGUIEvent( + { + type: EventType.CUSTOM, + name: "langflow.human_input_required", + value: { request_id: "node:job-1", kind: "node_input" }, + } as unknown as BaseEvent, + ctx, + ); + + expect(terminal).toBe(false); // pause is NOT terminal + expect(calls).toEqual(["humanInput:node:job-1"]); + }); +}); diff --git a/src/frontend/src/controllers/API/agui/run-agent.ts b/src/frontend/src/controllers/API/agui/run-agent.ts index 0f6e573053..b1e9c12a56 100644 --- a/src/frontend/src/controllers/API/agui/run-agent.ts +++ b/src/frontend/src/controllers/API/agui/run-agent.ts @@ -189,3 +189,36 @@ export function createWorkflowAgent( }; return agent; } + +/** + * Agent that reattaches to a background run's durable event stream via + * `GET /api/v2/workflows/{job_id}/events` (LE-1442/1450). HITL runs are submitted + * in background mode (a separate JSON POST) and consumed here; `lastEventId` is the + * SSE `Last-Event-ID` so a resume continues gap-free past the suspend boundary. + */ +export function createEventsAgent(opts: { + jobId: string; + lastEventId?: string; + headers?: Record; +}): HttpAgent { + const config: HttpAgentConfig = { + url: `${WORKFLOWS_ENDPOINT}/${encodeURIComponent(opts.jobId)}/events`, + headers: opts.headers, + }; + const agent = new HttpAgent(config); + ( + agent as unknown as { requestInit: (input: RunAgentInput) => RequestInit } + ).requestInit = function requestInit(_input: RunAgentInput): RequestInit { + const headers: Record = { + ...agent.headers, + Accept: "text/event-stream", + }; + if (opts.lastEventId) headers["Last-Event-ID"] = opts.lastEventId; + return { + method: "GET", + headers, + signal: agent.abortController.signal, + }; + }; + return agent; +} diff --git a/src/frontend/src/controllers/API/agui/run-flow-bridge.ts b/src/frontend/src/controllers/API/agui/run-flow-bridge.ts index 74ed69b1f9..03b84c21d2 100644 --- a/src/frontend/src/controllers/API/agui/run-flow-bridge.ts +++ b/src/frontend/src/controllers/API/agui/run-flow-bridge.ts @@ -10,18 +10,25 @@ import { type BaseEvent, EventType } from "@ag-ui/client"; import { handleMessageEvent } from "@/components/core/playgroundComponent/chat-view/utils/message-event-handler"; +import { updateMessage } from "@/components/core/playgroundComponent/chat-view/utils/message-utils"; +import { queryClient } from "@/contexts"; import { BuildStatus } from "@/constants/enums"; import useAlertStore from "@/stores/alertStore"; import useFlowStore from "@/stores/flowStore"; +import { useMessagesStore } from "@/stores/messagesStore"; import type { ChatInputType, ChatOutputType, VertexBuildTypeAPI, VertexDataTypeAPI, } from "@/types/api"; +import type { ContentBlock, InteractiveContent } from "@/types/chat"; +import type { Message } from "@/types/messages"; +import { api } from "../api"; import { buildWorkflowRunRequest, createWorkflowAgent, + WORKFLOWS_ENDPOINT, WORKFLOWS_PUBLIC_ENDPOINT, type WorkflowRunOptions, } from "./run-agent"; @@ -53,6 +60,7 @@ export interface BridgeContext { setRunId: (runId: string) => void; applyDelta: (ops: JsonPatchOp[]) => void; handleCustomEvent: (eventType: string, data: unknown) => void; + onHumanInput: (payload: unknown) => void; onFinished: () => void; onError: (message: string) => void; } @@ -83,6 +91,9 @@ export function handleAGUIEvent(event: BaseEvent, ctx: BridgeContext): boolean { }; if (custom.name === "langflow.event" && custom.value?.event_type) { ctx.handleCustomEvent(custom.value.event_type, custom.value.data); + } else if (custom.name === "langflow.human_input_required") { + // Non-terminal: surface the card; the SSE ends at suspend and reattaches on resume. + ctx.onHumanInput(custom.value); } return false; } @@ -239,6 +250,11 @@ export async function runFlowAGUI( }, applyDelta: (ops) => applyStateDelta(ops, runId, touchedNodeIds), handleCustomEvent: (eventType, data) => handleMessageEvent(eventType, data), + onHumanInput: () => { + // Stream mode can't durably suspend; HITL runs go through the background + // path (runFlowHITL). Mark awaiting-input so the UI reflects the pause. + flowStore.setAwaitingInput(true); + }, onFinished: () => { terminalEventSeen = true; flowStore.setBuildInfo({ success: true }); @@ -321,3 +337,289 @@ export async function runFlowAGUI( }); }); } + +/** Background-run body (mode="background"); HITL needs the durable substrate. */ +function buildBackgroundRunRequest(opts: WorkflowRunOptions) { + const body: Record = { + flow_id: opts.flowId, + input_value: opts.message ?? "", + mode: "background", + stream_protocol: "agui", + }; + if (opts.threadId) body.session_id = opts.threadId; + if (opts.tweaks) body.tweaks = opts.tweaks; + if (opts.startComponentId) body.start_component_id = opts.startComponentId; + if (opts.stopComponentId) body.stop_component_id = opts.stopComponentId; + if (opts.flowData) body.data = opts.flowData; + if (opts.files && opts.files.length > 0) body.files = opts.files; + return body; +} + +function toInteractiveContent( + payload: Record, + jobId: string, +): InteractiveContent { + const allowed = (payload.allowed_decisions as string[]) ?? []; + return { + type: "human_input", + kind: (payload.kind as InteractiveContent["kind"]) ?? "node_input", + request_id: String(payload.request_id ?? ""), + prompt: payload.prompt as string | undefined, + options: (payload.options as InteractiveContent["options"]) ?? [], + schema: payload.schema as InteractiveContent["schema"], + allowed_decisions: allowed, + job_id: jobId, + }; +} + +/** + * Reattach context per pause, keyed by request id. The interactive card resumes the + * run itself (no prop-drilling through the chat render chain), so it needs the exact + * run opts to stream the continued run back into the right session. + */ +const resumeRegistry = new Map< + string, + { jobId: string; opts: WorkflowRunOptions } +>(); + +export function getResumeContext( + requestId: string, +): { jobId: string; opts: WorkflowRunOptions } | undefined { + return resumeRegistry.get(requestId); +} + +function _cardSubmittedAction(messageId: string): string | undefined { + for (const query of queryClient.getQueryCache().getAll()) { + const key = query.queryKey; + if (!Array.isArray(key) || key[0] !== "useGetMessagesQuery") continue; + const messages = query.state.data as Message[] | undefined; + const msg = Array.isArray(messages) + ? messages.find((m) => m.id === messageId) + : undefined; + const content = msg?.content_blocks?.[0]?.contents?.find( + (c) => c?.type === "human_input", + ) as { submitted_action?: string } | undefined; + if (content?.submitted_action) return content.submitted_action; + } + return undefined; +} + +function cardAlreadyAnswered(messageId: string): boolean { + return _cardSubmittedAction(messageId) !== undefined; +} + +/** + * Stamp the chosen action onto the card message in the react-query cache so the + * selection is derived from a stable source — local React state is lost when the + * resume reattach replays the stream and re-renders the message list. + */ +export function markHumanInputSubmitted( + requestId: string, + actionId: string, +): void { + const messageId = `human-input-${requestId}`; + for (const query of queryClient.getQueryCache().getAll()) { + const key = query.queryKey; + if (!Array.isArray(key) || key[0] !== "useGetMessagesQuery") continue; + const messages = query.state.data as Message[] | undefined; + if (!Array.isArray(messages) || !messages.some((m) => m.id === messageId)) { + continue; + } + queryClient.setQueryData(key, (old: Message[] = []) => + old.map((m) => + m.id === messageId + ? { + ...m, + content_blocks: (m.content_blocks ?? []).map((block) => ({ + ...block, + contents: (block.contents ?? []).map((c) => + c?.type === "human_input" + ? { ...c, submitted_action: actionId } + : c, + ), + })), + } + : m, + ), + ); + } +} + +/** Render the pause as an interactive card in the chat and flag awaiting-input. */ +function injectHumanInputCard( + payload: Record, + jobId: string, + opts: WorkflowRunOptions, +): void { + const content = toInteractiveContent(payload, jobId); + resumeRegistry.set(content.request_id, { jobId, opts }); + const messageId = `human-input-${content.request_id}`; + // A resume reattach replays the pause event; if the user already answered (the + // card carries submitted_action), re-injecting would clobber that choice. Skip. + if (cardAlreadyAnswered(messageId)) { + useFlowStore.getState().setAwaitingInput(true); + return; + } + const block: ContentBlock = { + title: "Human input required", + contents: [content], + allow_markdown: true, + component: "HumanInput", + }; + const message: Message = { + flow_id: opts.flowId, + text: "", + sender: "Machine", + sender_name: "AI", + session_id: opts.threadId ?? opts.flowId, + timestamp: new Date().toISOString(), + files: [], + id: messageId, + edit: false, + background_color: "", + text_color: "", + content_blocks: [block], + }; + // The new playground reads the react-query messages cache; the legacy IOModal + // playground reads useMessagesStore. Write both so the card renders either way. + updateMessage(message); + useMessagesStore.getState().addMessage(message); + useFlowStore.getState().setAwaitingInput(true); +} + +/** + * Consume a background run's durable event stream (`GET /{job_id}/events`) and fold + * events into the flow store. Resolves when the run ends, suspends for human input, + * or the request errors. Reused for the initial run and for resume reattach. + * + * Uses a direct fetch + SSE parser rather than the AG-UI ``HttpAgent``: a paused run + * closes the stream WITHOUT a terminal ``RUN_FINISHED``, which the agent's protocol + * verifier treats as a violation and can swallow the pause event we depend on. + */ +export async function consumeBackgroundEvents( + jobId: string, + opts: WorkflowRunOptions & { signal?: AbortSignal }, + lastEventId?: string, +): Promise { + const flowStore = useFlowStore.getState(); + const setErrorData = useAlertStore.getState().setErrorData; + const touchedNodeIds = new Set(); + let runId = jobId; + let suspended = false; + + const ctx: BridgeContext = { + setRunId: (r) => { + runId = r; + }, + applyDelta: (ops) => applyStateDelta(ops, runId, touchedNodeIds), + handleCustomEvent: (eventType, data) => handleMessageEvent(eventType, data), + onHumanInput: (payload) => { + suspended = true; + injectHumanInputCard(payload as Record, jobId, opts); + }, + onFinished: () => flowStore.setBuildInfo({ success: true }), + onError: (message) => { + flowStore.setBuildInfo({ error: [message], success: false }); + setErrorData({ title: "Workflow run failed", list: [message] }); + }, + }; + + const finish = () => { + flowStore.updateEdgesRunningByNodes([...touchedNodeIds], false); + // A suspended run stops building and parks awaiting input; the spinner reads + // isBuilding, so it must clear on suspend too — only awaitingInput differs. + flowStore.setIsBuilding(false); + flowStore.setAwaitingInput(suspended); + flowStore.revertBuiltStatusFromBuilding(); + }; + + const headers: Record = { Accept: "text/event-stream" }; + if (lastEventId) headers["Last-Event-ID"] = lastEventId; + const url = `${WORKFLOWS_ENDPOINT}/${encodeURIComponent(jobId)}/events`; + + try { + const response = await fetch(url, { + method: "GET", + headers, + credentials: "same-origin", + signal: opts.signal, + }); + if (!response.ok || !response.body) { + ctx.onError(`Events stream failed with status ${response.status}`); + finish(); + return; + } + const reader = response.body.getReader(); + const decoder = new TextDecoder(); + let buffer = ""; + let terminal = false; + + const dispatch = (frame: string): boolean => { + const dataLine = frame + .split("\n") + .find((line) => line.startsWith("data:")); + if (!dataLine) return false; + let event: BaseEvent; + try { + event = JSON.parse(dataLine.slice(5).trim()) as BaseEvent; + } catch { + return false; + } + if (handleAGUIEvent(event, ctx)) { + // A terminal event supersedes a replayed pause (resume reattach re-emits it). + suspended = false; + return true; + } + return false; + }; + + while (!terminal) { + const { value, done } = await reader.read(); + if (done) { + // The stream can close right after the pause event without a trailing + // blank line; flush whatever is buffered so that last frame isn't lost. + if (buffer.trim()) dispatch(buffer); + break; + } + buffer += decoder.decode(value, { stream: true }); + let boundary = buffer.indexOf("\n\n"); + while (boundary !== -1) { + const frame = buffer.slice(0, boundary); + buffer = buffer.slice(boundary + 2); + boundary = buffer.indexOf("\n\n"); + if (dispatch(frame)) { + terminal = true; + break; + } + } + } + finish(); + } catch (err) { + if ((err as Error)?.name !== "AbortError") { + ctx.onError((err as Error)?.message ?? "Events stream error"); + } + finish(); + } +} + +/** + * Run a flow that contains a Human Input node through the durable background path: + * submit in background mode, then consume the run's event stream. The pause renders + * an interactive card; resume reattaches via {@link consumeBackgroundEvents}. + */ +export async function runFlowHITL( + opts: WorkflowRunOptions & { signal?: AbortSignal }, +): Promise { + const body = buildBackgroundRunRequest(opts); + const { data } = await api.post(WORKFLOWS_ENDPOINT, body); + const jobId: string | undefined = data?.job_id; + if (!jobId) { + useFlowStore.getState().setBuildInfo({ + error: ["Background run did not return a job id"], + success: false, + }); + useFlowStore.getState().setIsBuilding(false); + return; + } + await consumeBackgroundEvents(jobId, opts); +} diff --git a/src/frontend/src/controllers/API/queries/workflows/use-resume-workflow.ts b/src/frontend/src/controllers/API/queries/workflows/use-resume-workflow.ts new file mode 100644 index 0000000000..e288815f1b --- /dev/null +++ b/src/frontend/src/controllers/API/queries/workflows/use-resume-workflow.ts @@ -0,0 +1,43 @@ +import type { useMutationFunctionType } from "@/types/api"; +import { WORKFLOWS_ENDPOINT } from "../../agui/run-agent"; +import { api } from "../../api"; +import { UseRequestProcessor } from "../../services/request-processor"; + +export interface ResumeWorkflowPayload { + jobId: string; + requestId: string; + decision: { action_id: string; values: Record }; +} + +export interface ResumeWorkflowResponse { + job_id: string; + status: string; + message?: string; +} + +/** + * Resume a SUSPENDED human-in-the-loop run with the human's decision. + * + * POSTs the decision to `/api/v2/workflows/{job_id}/resume` (LE-1450). The caller + * re-attaches a fresh `GET /{job_id}/events` from the last seen event id so the + * continued run streams gap-free. + */ +export const useResumeWorkflow: useMutationFunctionType< + undefined, + ResumeWorkflowPayload, + ResumeWorkflowResponse +> = (options?) => { + const { mutate } = UseRequestProcessor(); + + const resumeFn = async ( + payload: ResumeWorkflowPayload, + ): Promise => { + const res = await api.post( + `${WORKFLOWS_ENDPOINT}/${encodeURIComponent(payload.jobId)}/resume`, + { request_id: payload.requestId, decision: payload.decision }, + ); + return res.data; + }; + + return mutate(["useResumeWorkflow"], resumeFn, options); +}; diff --git a/src/frontend/src/stores/flowStore.ts b/src/frontend/src/stores/flowStore.ts index 4e82bf5bbd..e9b6814339 100644 --- a/src/frontend/src/stores/flowStore.ts +++ b/src/frontend/src/stores/flowStore.ts @@ -10,7 +10,10 @@ import { cloneDeep } from "lodash"; import { create } from "zustand"; import { checkCodeValidity } from "@/CustomNodes/helpers/check-code-validity"; import { queryClient } from "@/contexts"; -import { runFlowAGUI } from "@/controllers/API/agui/run-flow-bridge"; +import { + runFlowAGUI, + runFlowHITL, +} from "@/controllers/API/agui/run-flow-bridge"; import { ENABLE_INSPECTION_PANEL } from "@/customization/feature-flags"; import { track, trackFlowBuild } from "@/customization/utils/analytics"; import { brokenEdgeMessage } from "@/utils/utils"; @@ -168,6 +171,7 @@ const useFlowStore = create((set, get) => ({ nodes: [], edges: [], isBuilding: false, + awaitingInput: false, buildStartTime: null, buildDuration: null, buildingFlowId: null, @@ -368,6 +372,9 @@ const useFlowStore = create((set, get) => ({ buildingSessionId: !isBuilding ? null : current.buildingSessionId, }); }, + setAwaitingInput: (awaitingInput) => { + set({ awaitingInput }); + }, setBuildStartTime: (time) => { set({ buildStartTime: time }); }, @@ -926,10 +933,8 @@ const useFlowStore = create((set, get) => ({ const buildController = new AbortController(); get().setBuildController(buildController); - // Always run through the v2 workflows endpoint. Current frontend nodes - // + edges are sent so unsaved tweaks (dropdowns, text inputs) run as - // the user sees them. - await runFlowAGUI({ + // A Human Input node can suspend mid-run — only the durable background path supports it. + const runArgs = { flowId: currentFlow!.id, message: input_value, threadId: session, @@ -938,7 +943,11 @@ const useFlowStore = create((set, get) => ({ flowData: { nodes: get().nodes, edges: get().edges }, files, signal: buildController.signal, - }); + }; + const hasHumanInput = get().nodes.some( + (node) => node.data?.type === "HumanInput", + ); + await (hasHumanInput ? runFlowHITL(runArgs) : runFlowAGUI(runArgs)); // Invalidate KB-related caches so any KnowledgeIngestion node that ran // inside this build surfaces its updated stats / runs the next time the diff --git a/src/frontend/src/types/chat/index.ts b/src/frontend/src/types/chat/index.ts index 81a88e5540..8eaa13dbbb 100644 --- a/src/frontend/src/types/chat/index.ts +++ b/src/frontend/src/types/chat/index.ts @@ -133,6 +133,32 @@ export interface ToolContent extends BaseContent { error?: JSONValue | string; } +// A human-in-the-loop pause: a decision the user must make before the run resumes. +export interface HumanInputOption { + action_id: string; + label?: string; +} + +export interface HumanInputFormField { + name: string; + type?: string; + required?: boolean; +} + +export interface InteractiveContent extends BaseContent { + type: "human_input"; + kind: "tool_approval" | "node_input"; + request_id: string; + prompt?: string; + options: HumanInputOption[]; + schema?: HumanInputFormField[]; + fields?: HumanInputFormField[]; + allowed_decisions: string[]; + // Correlation for the resume round-trip; stamped from the run that paused. + job_id?: string; + submitted_action?: string; +} + // Union type for all content types export type ContentType = | ErrorContent @@ -140,7 +166,8 @@ export type ContentType = | MediaContent | JSONContent | CodeContent - | ToolContent; + | ToolContent + | InteractiveContent; // Updated ContentBlock interface export interface ContentBlock { diff --git a/src/frontend/src/types/zustand/flow/index.ts b/src/frontend/src/types/zustand/flow/index.ts index 08928b463a..c5a6fce920 100644 --- a/src/frontend/src/types/zustand/flow/index.ts +++ b/src/frontend/src/types/zustand/flow/index.ts @@ -118,7 +118,9 @@ export type FlowStoreType = { buildingFlowId: string | null; buildingSessionId: string | null; isPending: boolean; + awaitingInput: boolean; setIsBuilding: (isBuilding: boolean) => void; + setAwaitingInput: (awaitingInput: boolean) => void; setBuildStartTime: (time: number) => void; setBuildDuration: (duration: number) => void; setBuildingSession: (flowId: string | null, sessionId: string | null) => void; diff --git a/src/lfx/src/lfx/graph/checkpoint/schema.py b/src/lfx/src/lfx/graph/checkpoint/schema.py index 1956db95ae..007368ed68 100644 --- a/src/lfx/src/lfx/graph/checkpoint/schema.py +++ b/src/lfx/src/lfx/graph/checkpoint/schema.py @@ -59,11 +59,16 @@ def serialize_value(value: Any) -> dict[str, Any] | None: if value is None or isinstance(value, (str, int, float, bool)): return {_WIRE_KIND: "raw", "value": value} if isinstance(value, BaseModel): + try: + dumped = value.model_dump(mode="json") + except Exception: # noqa: BLE001 + # A model with an opaque field (LLM client, model class) can't round-trip; degrade to None so the checkpoint stays writable. + return None return { _WIRE_KIND: "model", "module": type(value).__module__, "name": type(value).__qualname__, - "value": value.model_dump(mode="json"), + "value": dumped, } if isinstance(value, dict): if not all(isinstance(k, str) for k in value): diff --git a/src/lfx/src/lfx/schema/content_block.py b/src/lfx/src/lfx/schema/content_block.py index 6df7ef6e8b..bbe266cc41 100644 --- a/src/lfx/src/lfx/schema/content_block.py +++ b/src/lfx/src/lfx/schema/content_block.py @@ -3,7 +3,15 @@ from typing import Annotated from pydantic import BaseModel, Discriminator, Field, Tag, field_serializer, field_validator from typing_extensions import TypedDict -from .content_types import CodeContent, ErrorContent, JSONContent, MediaContent, TextContent, ToolContent +from .content_types import ( + CodeContent, + ErrorContent, + HumanInputContent, + JSONContent, + MediaContent, + TextContent, + ToolContent, +) def _get_type(d: dict | BaseModel) -> str | None: @@ -19,7 +27,8 @@ ContentType = Annotated[ | Annotated[TextContent, Tag("text")] | Annotated[MediaContent, Tag("media")] | Annotated[CodeContent, Tag("code")] - | Annotated[JSONContent, Tag("json")], + | Annotated[JSONContent, Tag("json")] + | Annotated[HumanInputContent, Tag("human_input")], Discriminator(_get_type), ] diff --git a/src/lfx/src/lfx/schema/content_types.py b/src/lfx/src/lfx/schema/content_types.py index 0682be4d88..0915e8001e 100644 --- a/src/lfx/src/lfx/schema/content_types.py +++ b/src/lfx/src/lfx/schema/content_types.py @@ -78,6 +78,24 @@ class CodeContent(BaseContent): title: str | None = None +class HumanInputContent(BaseContent): + """Content type for a human-in-the-loop pause persisted in the chat history. + + Carries the pending decision so the interactive card survives reload: the + request_id and job_id let the card resume the suspended run after an F5. + """ + + type: Literal["human_input"] = Field(default="human_input") + request_id: str + job_id: str | None = None + kind: str = "node_input" + prompt: str | None = None + options: list[dict[str, Any]] = Field(default_factory=list) + fields: list[dict[str, Any]] = Field(default_factory=list) + allowed_decisions: list[str] = Field(default_factory=list) + submitted_action: str | None = None + + class ToolContent(BaseContent): """Content type for tool start content.""" diff --git a/src/lfx/src/lfx/schema/workflow.py b/src/lfx/src/lfx/schema/workflow.py index b48c1561f4..22f95e8384 100644 --- a/src/lfx/src/lfx/schema/workflow.py +++ b/src/lfx/src/lfx/schema/workflow.py @@ -457,6 +457,21 @@ class WorkflowStopResponse(BaseModel): message: str | None = None +class WorkflowResumeRequest(BaseModel): + """Request schema for resuming a suspended (human-in-the-loop) workflow.""" + + request_id: str + decision: dict | None = None + + +class WorkflowResumeResponse(BaseModel): + """Response schema for resuming a workflow.""" + + job_id: JobId + status: str + message: str | None = None + + # OpenAPI response definitions WORKFLOW_EXECUTION_RESPONSES = { 200: { diff --git a/src/lfx/tests/unit/graph/checkpoint/_static_pauser.py b/src/lfx/tests/unit/graph/checkpoint/_static_pauser.py new file mode 100644 index 0000000000..e4b188dcbc --- /dev/null +++ b/src/lfx/tests/unit/graph/checkpoint/_static_pauser.py @@ -0,0 +1,42 @@ +"""A real, code-wireable component that pauses then resumes on an injected decision. + +HumanInput's branch outputs are dynamic and can't be wired via ``.set()`` in a +unit test, so this helper mirrors HumanInput's pause/resume CONTRACT with a single +static output: first run requests a pause; on resume it reads the decision the +build path injects into ``graph.human_input_decisions`` and carries it downstream. +This exercises the real resume_from_checkpoint + inject + un-build mechanics that +``build.py``'s resume branch relies on, with a graph that wires in code. +""" + +from __future__ import annotations + +from lfx.custom import Component +from lfx.io import MessageTextInput, Output +from lfx.schema.message import Message + +HUMAN_INPUT_REQUIRED = "human_input_required" + + +class StaticPauser(Component): + display_name = "Static Pauser" + name = "StaticPauser" + + inputs = [MessageTextInput(name="input_value", display_name="In")] + outputs = [Output(display_name="Out", name="out", method="run_it")] + + def _request_id(self) -> str: + return f"{self._id}:{self.graph.run_id}" + + def _decision(self) -> dict | None: + decisions = getattr(self.graph, "human_input_decisions", None) + return decisions.get(self._request_id()) if isinstance(decisions, dict) else None + + def run_it(self) -> Message: + decision = self._decision() + if decision is None: + self.graph.request_pause( + reason=HUMAN_INPUT_REQUIRED, + data={"request_id": self._request_id(), "kind": "node_input"}, + ) + return Message(text="") + return Message(text=str(decision.get("action_id", ""))) diff --git a/src/lfx/tests/unit/graph/checkpoint/test_resume_round_trip.py b/src/lfx/tests/unit/graph/checkpoint/test_resume_round_trip.py new file mode 100644 index 0000000000..6f078af8ce --- /dev/null +++ b/src/lfx/tests/unit/graph/checkpoint/test_resume_round_trip.py @@ -0,0 +1,89 @@ +"""End-to-end graph resume round-trip (LE-1449 / LE-1446 mechanics). + +A real ChatInput -> StaticPauser -> ChatOutput graph: first run suspends with a +node_input request and a durable checkpoint; resume hydrates the checkpoint, +injects the decision, un-builds the paused node, and runs to completion without +re-executing already-built vertices. This proves the exact resume-from-checkpoint ++ inject + un-build sequence that build.py's resume branch performs. +""" + +from __future__ import annotations + +import sys +from pathlib import Path + +import pytest +from lfx.components.input_output import ChatInput, ChatOutput +from lfx.graph import Graph +from lfx.graph.checkpoint.store import InMemoryCheckpointStore +from lfx.graph.exceptions import GraphPausedException + +sys.path.insert(0, str(Path(__file__).parent)) +from _static_pauser import StaticPauser + + +def _graph(store, *, job_id="job-1"): + chat_input = ChatInput(_id="chat_input", input_value="hello") + chat_input.set(should_store_message=False) + pauser = StaticPauser(_id="pauser") + pauser.set(input_value=chat_input.message_response) + chat_output = ChatOutput(_id="chat_output") + chat_output.set(input_value=pauser.run_it, should_store_message=False) + graph = Graph(chat_input, chat_output) + graph.session_id = "sess-1" + graph.set_run_id(job_id) + graph.job_id = job_id + graph.checkpointing_enabled = True + graph.checkpoint_store = store + return graph + + +async def test_first_run_suspends_with_durable_checkpoint(): + store = InMemoryCheckpointStore() + graph = _graph(store) + + with pytest.raises(GraphPausedException) as excinfo: + await graph.process(fallback_to_env_vars=False) + + assert excinfo.value.reason == "human_input_required" + assert excinfo.value.data["request_id"] == "pauser:job-1" + assert graph.get_vertex("chat_input").built is True + assert graph.get_vertex("chat_output").built is False + assert await store.load_by_run_id("job-1") is not None + + +async def test_resume_injects_decision_runs_to_completion_without_reexec(): + store = InMemoryCheckpointStore() + first = _graph(store) + with pytest.raises(GraphPausedException): + await first.process(fallback_to_env_vars=False) + + # The build.py resume branch performs exactly these four steps: + checkpoint = await store.load_by_run_id("job-1") + resumed = Graph.resume_from_checkpoint(checkpoint, checkpoint_store=store) + resumed.human_input_decisions = {"pauser:job-1": {"action_id": "approve", "values": {}}} + resumed.get_vertex("pauser").built = False + + assert resumed.resume_first_layer() == ["pauser"] + assert resumed.get_vertex("chat_input").built is True # not re-run + + await resumed.process(fallback_to_env_vars=False) + + assert resumed.get_vertex("pauser").built is True + assert resumed.get_vertex("chat_output").built is True + pauser_result = resumed.get_vertex("pauser").results + assert pauser_result["out"].text == "approve" # the injected decision flowed downstream + + +async def test_resume_without_decision_suspends_again(): + store = InMemoryCheckpointStore() + first = _graph(store) + with pytest.raises(GraphPausedException): + await first.process(fallback_to_env_vars=False) + + checkpoint = await store.load_by_run_id("job-1") + resumed = Graph.resume_from_checkpoint(checkpoint, checkpoint_store=store) + resumed.get_vertex("pauser").built = False # re-run, but no decision injected + + with pytest.raises(GraphPausedException): + await resumed.process(fallback_to_env_vars=False) diff --git a/src/lfx/tests/unit/graph/checkpoint/test_schema.py b/src/lfx/tests/unit/graph/checkpoint/test_schema.py index 82af769f89..5b273679c7 100644 --- a/src/lfx/tests/unit/graph/checkpoint/test_schema.py +++ b/src/lfx/tests/unit/graph/checkpoint/test_schema.py @@ -109,3 +109,18 @@ def test_serialize_value_returns_none_for_opaque_objects(): assert serialize_value(Opaque()) is None assert serialize_value(lambda: 1) is None + + +def test_serialize_value_degrades_model_with_unserializable_field(): + """A model holding an opaque field (e.g. an LLM client / model class) must not raise. + + Reproduces the HITL-with-Agent crash: pausing serialized the agent's state and a + nested model class blew up ``model_dump(mode="json")``. + """ + from pydantic import BaseModel, ConfigDict + + class Holder(BaseModel): + model_config = ConfigDict(arbitrary_types_allowed=True) + opaque: type = BaseModel + + assert serialize_value(Holder()) is None