refactor(lfx): extract v2 workflow contract layer into lfx.workflow

Moves the protocol-agnostic pieces of the v2 workflows API out of the langflow
backend into lfx so both the backend and `lfx serve` can share one contract.
First step toward giving lfx (the production runtime) the v2 workflows API.

- Move api/v2/adapters/, agui_translator.py, and converters.py to lfx/workflow/.
  They depend only on lfx.schema.workflow and ag_ui (already an lfx dep), so lfx
  carries the contract with zero langflow imports.
- Decouple the one langflow reference: converters typed run_response against
  langflow.api.v1.schemas.RunResponse (TYPE_CHECKING only). Replaced with a local
  RunResponseLike Protocol (outputs + session_id), the only attributes used.
- Repoint the six backend v2 workflow modules to import from lfx.workflow.
- Move the five protocol-agnostic contract tests into src/lfx/tests/unit/workflow/
  (run in the lfx-only env). test_output_event_parity and test_workflow_agui stay
  in langflow (they need langflow.api.build) with repointed imports.

Coverage unchanged: 201 contract tests pass in the lfx-only env, 191 backend v2
tests pass; 392 total, same as before the move.
This commit is contained in:
ogabrielluiz
2026-06-23 12:54:54 -03:00
parent a21d0d9eab
commit 7116ea1842
21 changed files with 108 additions and 76 deletions

View File

@ -43,17 +43,17 @@ from lfx.schema.workflow import (
WorkflowStopResponse,
)
from lfx.services.deps import injectable_session_scope_readonly
from pydantic_core import ValidationError as PydanticValidationError
from sqlalchemy.exc import OperationalError
from langflow.api.v2.adapters import (
from lfx.workflow.adapters import (
STREAM_ADAPTERS,
StreamAdapterContext,
UnknownStreamProtocolError,
available_protocols,
get_stream_adapter,
)
from langflow.api.v2.converters import parse_workflow_run_request
from lfx.workflow.converters import parse_workflow_run_request
from pydantic_core import ValidationError as PydanticValidationError
from sqlalchemy.exc import OperationalError
from langflow.api.v2.workflow_background import (
_BACKGROUND_RUNS,
_cancel_workflow_queue_job,

View File

@ -19,14 +19,14 @@ from fastapi import BackgroundTasks, Request
from fastapi.sse import format_sse_event
from lfx.log.logger import logger
from lfx.schema.workflow import JobId, JobStatus, WorkflowJobResponse
from langflow.api.v2.adapters import (
from lfx.workflow.adapters import (
StreamAdapterContext,
StreamEvent,
UnknownStreamProtocolError,
get_stream_adapter,
)
from langflow.api.v2.converters import ParsedWorkflowRun
from lfx.workflow.converters import ParsedWorkflowRun
from langflow.api.v2.workflow_execution import _stream_event_frames
from langflow.exceptions.api import (
WorkflowQueueFullError,

View File

@ -34,11 +34,11 @@ from lfx.graph.graph.base import Graph
from lfx.log.logger import logger
from lfx.schema.schema import InputValueRequest
from lfx.schema.workflow import WorkflowExecutionResponse
from lfx.workflow.adapters import StreamAdapter, StreamEvent
from lfx.workflow.converters import ParsedWorkflowRun, create_error_response, run_response_to_workflow_response
from langflow.api.utils import extract_global_variables_from_headers
from langflow.api.v1.schemas import FlowDataRequest, RunResponse
from langflow.api.v2.adapters import StreamAdapter, StreamEvent
from langflow.api.v2.converters import ParsedWorkflowRun, create_error_response, run_response_to_workflow_response
from langflow.api.v2.workflow_validation import _validate_output_ids
from langflow.exceptions.api import WorkflowTimeoutError, WorkflowValidationError
from langflow.processing.process import process_tweaks, run_graph_internal

View File

@ -40,6 +40,14 @@ from lfx.utils.flow_validation import (
validate_flow_for_current_settings,
validate_public_flow_no_code_execution,
)
from lfx.workflow.adapters import (
STREAM_ADAPTERS,
StreamAdapterContext,
UnknownStreamProtocolError,
available_protocols,
get_stream_adapter,
)
from lfx.workflow.converters import ParsedWorkflowRun
from limits import parse
from langflow.api.utils.flow_utils import (
@ -47,14 +55,6 @@ from langflow.api.utils.flow_utils import (
validate_public_files,
verify_public_flow_and_get_user,
)
from langflow.api.v2.adapters import (
STREAM_ADAPTERS,
StreamAdapterContext,
UnknownStreamProtocolError,
available_protocols,
get_stream_adapter,
)
from langflow.api.v2.converters import ParsedWorkflowRun
from langflow.services.auth.utils import get_current_user_optional
from langflow.services.database.models.flow.model import Flow
from langflow.services.database.models.user.model import User, UserRead

View File

@ -10,9 +10,9 @@ from typing import TYPE_CHECKING
from lfx.graph.graph.base import Graph
from lfx.graph.schema import ResultData, RunOutputs
from lfx.workflow.converters import run_response_to_workflow_response
from langflow.api.v1.schemas import RunResponse
from langflow.api.v2.converters import run_response_to_workflow_response
from langflow.services.database.models.vertex_builds.crud import get_vertex_builds_by_job_id
if TYPE_CHECKING:

View File

@ -10,8 +10,8 @@ from __future__ import annotations
from fastapi import HTTPException, status
from lfx.utils.flow_validation import CustomComponentValidationError, validate_flow_for_current_settings
from lfx.workflow.converters import ParsedWorkflowRun
from langflow.api.v2.converters import ParsedWorkflowRun
from langflow.services.authorization.fetch import deny_to_404
from langflow.services.database.models.flow.model import FlowRead
from langflow.services.database.models.user.model import UserRead

View File

@ -15,9 +15,9 @@ import json
from unittest.mock import Mock
from langflow.api.build import _output_meta_for_vertex
from langflow.api.v2.adapters import StreamAdapterContext
from langflow.api.v2.adapters.langflow import LangflowAdapter
from langflow.api.v2.converters import build_component_output, resolve_output_type
from lfx.workflow.adapters import StreamAdapterContext
from lfx.workflow.adapters.langflow import LangflowAdapter
from lfx.workflow.converters import build_component_output, resolve_output_type
def _adapter() -> LangflowAdapter:

View File

@ -399,8 +399,8 @@ class TestAGUIStreaming:
async def test_stream_event_handoff_overflow_emits_error(self, monkeypatch: pytest.MonkeyPatch):
"""EventManager put_nowait must not silently drop frames when the stream buffer fills."""
from langflow.api.v2 import workflow_execution as wf_exec
from langflow.api.v2.adapters import StreamEvent
from langflow.api.v2.converters import ParsedWorkflowRun
from lfx.workflow.adapters import StreamEvent
from lfx.workflow.converters import ParsedWorkflowRun
seen_maxsize: list[int] = []
@ -451,8 +451,8 @@ class TestAGUIStreaming:
):
"""A run that exceeds the wall-clock ceiling ends in a sanitized terminal error, not a hang."""
from langflow.api.v2 import workflow_execution as wf_exec
from langflow.api.v2.adapters import StreamAdapterContext, get_stream_adapter
from langflow.api.v2.converters import ParsedWorkflowRun
from lfx.workflow.adapters import StreamAdapterContext, get_stream_adapter
from lfx.workflow.converters import ParsedWorkflowRun
async def hanging_generate_flow_events(**_kwargs):
await asyncio.sleep(5)
@ -490,8 +490,8 @@ class TestAGUIStreaming:
):
"""A producer that reports on_error then raises must not triple-emit terminal errors."""
from langflow.api.v2 import workflow_execution as wf_exec
from langflow.api.v2.adapters import StreamAdapterContext, get_stream_adapter
from langflow.api.v2.converters import ParsedWorkflowRun
from lfx.workflow.adapters import StreamAdapterContext, get_stream_adapter
from lfx.workflow.converters import ParsedWorkflowRun
async def fake_generate_flow_events(**kwargs):
kwargs["event_manager"].on_error(data={"error": "inner"})
@ -528,8 +528,8 @@ class TestAGUIStreaming:
):
"""A producer that raises before on_error still emits one terminal error."""
from langflow.api.v2 import workflow_execution as wf_exec
from langflow.api.v2.adapters import StreamAdapterContext, get_stream_adapter
from langflow.api.v2.converters import ParsedWorkflowRun
from lfx.workflow.adapters import StreamAdapterContext, get_stream_adapter
from lfx.workflow.converters import ParsedWorkflowRun
async def fake_generate_flow_events(**_kwargs):
message = "early boom"
@ -562,7 +562,7 @@ class TestAGUIStreaming:
async def test_agui_stream_emits_end_side_channel_for_build_duration(self, monkeypatch: pytest.MonkeyPatch):
"""The AG-UI stream must preserve v1 end payloads for chat build-duration persistence."""
from langflow.api.v2 import workflow_execution as wf_exec
from langflow.api.v2.converters import ParsedWorkflowRun
from lfx.workflow.converters import ParsedWorkflowRun
async def fake_generate_flow_events(**kwargs):
event_queue = kwargs["event_manager"].queue
@ -1008,7 +1008,7 @@ class TestAGUIBackgroundJobStatus:
async def test_background_buffer_binds_build_rows_to_returned_job_id(self, monkeypatch: pytest.MonkeyPatch):
"""Status reconstruction needs vertex_build rows logged under the public background job id."""
from langflow.api.v2 import workflow_background as wf_bg
from langflow.api.v2.converters import ParsedWorkflowRun
from lfx.workflow.converters import ParsedWorkflowRun
job_id = uuid4()
captured: dict = {}
@ -1044,7 +1044,7 @@ class TestAGUIBackgroundJobStatus:
"""A background job should leave QUEUED once its buffer task starts executing."""
from langflow.api.v2 import workflow as workflow_module
from langflow.api.v2 import workflow_background as wf_bg
from langflow.api.v2.converters import ParsedWorkflowRun
from lfx.workflow.converters import ParsedWorkflowRun
job_id = uuid4()
updates: list[tuple[object, object, bool]] = []
@ -1086,7 +1086,7 @@ class TestAGUIBackgroundJobStatus:
"""The owner task must append cancellation before marking replay done."""
from langflow.api.v2 import workflow_background as wf_bg
from langflow.api.v2 import workflow_execution as wf_exec
from langflow.api.v2.converters import ParsedWorkflowRun
from lfx.workflow.converters import ParsedWorkflowRun
started = asyncio.Event()
@ -1151,7 +1151,7 @@ class TestAGUIBackgroundJobStatus:
"""Cancellation framing must stay protocol-native outside AG-UI too."""
from langflow.api.v2 import workflow_background as wf_bg
from langflow.api.v2 import workflow_execution as wf_exec
from langflow.api.v2.converters import ParsedWorkflowRun
from lfx.workflow.converters import ParsedWorkflowRun
started = asyncio.Event()
@ -1209,7 +1209,7 @@ class TestAGUIBackgroundJobStatus:
"""The stop fallback must not wake replay readers before cancellation is buffered."""
from langflow.api.v2 import workflow_background as wf_bg
from langflow.api.v2 import workflow_execution as wf_exec
from langflow.api.v2.converters import ParsedWorkflowRun
from lfx.workflow.converters import ParsedWorkflowRun
job_id = str(uuid4())
started = asyncio.Event()
@ -2256,13 +2256,13 @@ class TestBufferBackgroundRunUnknownProtocolGuard:
from uuid import uuid4 as _uuid4
from langflow.api.v2 import workflow_background as wf_bg
from langflow.api.v2.adapters import STREAM_ADAPTERS as _REGISTRY
from langflow.api.v2.converters import ParsedWorkflowRun
from langflow.services.database.models.flow.model import Flow, FlowRead
from langflow.services.database.models.jobs.model import Job, JobStatus
from langflow.services.database.models.user.model import User as _User
from langflow.services.database.models.user.model import UserRead
from langflow.services.deps import get_job_service
from lfx.workflow.adapters import STREAM_ADAPTERS as _REGISTRY
from lfx.workflow.converters import ParsedWorkflowRun
# Real Job row so update_job_status can flip it.
job_id = _uuid4()
@ -2311,12 +2311,12 @@ class TestBufferBackgroundRunUnknownProtocolGuard:
from uuid import uuid4 as _uuid4
from langflow.api.v2 import workflow_background as wf_bg
from langflow.api.v2.adapters import STREAM_ADAPTERS as _REGISTRY
from langflow.api.v2.converters import ParsedWorkflowRun
from langflow.services.database.models.flow.model import Flow, FlowRead
from langflow.services.database.models.user.model import User as _User
from langflow.services.database.models.user.model import UserRead
from langflow.services.deps import get_job_service
from lfx.workflow.adapters import STREAM_ADAPTERS as _REGISTRY
from lfx.workflow.converters import ParsedWorkflowRun
job_id = _uuid4()
await get_job_service().create_job(
@ -2360,7 +2360,7 @@ class TestExecuteWorkflowBackgroundQueueOwnership:
so registering them would let the polling watchdog reclaim long runs.
"""
from langflow.api.v2 import workflow_background as wf_bg
from langflow.api.v2.converters import ParsedWorkflowRun
from lfx.workflow.converters import ParsedWorkflowRun
job_id = uuid4()
current_user_id = uuid4()

View File

@ -0,0 +1,14 @@
"""V2 workflow contract layer shared by the langflow backend and ``lfx serve``.
This package holds the protocol-agnostic pieces of the V2 workflow API:
- ``adapters``: the ``StreamAdapter`` protocol, the ``langflow``/``agui`` SSE
adapters, and the registry (``get_stream_adapter`` / ``available_protocols``).
- ``agui_translator``: translates EventManager events into AG-UI events.
- ``converters``: parses ``WorkflowRunRequest`` and builds the structured
``WorkflowExecutionResponse``.
It depends only on ``lfx.schema.workflow`` and ``ag_ui`` (no langflow imports),
so both runtimes consume one contract. The langflow backend layers its stateful
"vehicle" (database, job queue, auth) on top.
"""

View File

@ -11,12 +11,16 @@ Adding a new protocol is one new module under ``adapters/`` plus one
from __future__ import annotations
from collections.abc import Callable, Iterable
from collections.abc import Callable
from dataclasses import dataclass
from typing import Any, ClassVar, Protocol, runtime_checkable
from typing import TYPE_CHECKING, Protocol, runtime_checkable
from pydantic import BaseModel
if TYPE_CHECKING:
from collections.abc import Iterable
from typing import Any, ClassVar
@dataclass(frozen=True)
class StreamEvent:
@ -117,5 +121,5 @@ def available_protocols() -> list[str]:
# Built-in adapter registrations happen here, after the registry is defined.
# Import for side-effect: each module calls ``register_stream_adapter``.
from langflow.api.v2.adapters import agui as _agui # noqa: E402, F401
from langflow.api.v2.adapters import langflow as _langflow # noqa: E402, F401
from lfx.workflow.adapters import agui as _agui # noqa: E402, F401
from lfx.workflow.adapters import langflow as _langflow # noqa: E402, F401

View File

@ -7,17 +7,19 @@ Per-run state lives on the wrapped translator instance.
from __future__ import annotations
from collections.abc import Iterable
from typing import Any, ClassVar
from typing import TYPE_CHECKING, Any, ClassVar
from ag_ui.core import BaseEvent
from langflow.api.v2.adapters import (
from lfx.workflow.adapters import (
StreamAdapterContext,
StreamEvent,
register_stream_adapter,
)
from langflow.api.v2.agui_translator import AGUITranslator
from lfx.workflow.agui_translator import AGUITranslator
if TYPE_CHECKING:
from collections.abc import Iterable
from ag_ui.core import BaseEvent
def _to_stream_event(event: BaseEvent) -> StreamEvent:

View File

@ -8,17 +8,18 @@ clients (curl users, the v1 frontend) can read it without learning anything new.
from __future__ import annotations
import json
from collections.abc import Iterable
from typing import Any, ClassVar
from typing import TYPE_CHECKING, Any, ClassVar
from lfx.schema.workflow import OutputEvent
from langflow.api.v2.adapters import (
from lfx.workflow.adapters import (
StreamAdapterContext,
StreamEvent,
register_stream_adapter,
)
from langflow.api.v2.converters import build_component_output, resolve_output_type
from lfx.workflow.converters import build_component_output, resolve_output_type
if TYPE_CHECKING:
from collections.abc import Iterable
class LangflowAdapter:

View File

@ -24,7 +24,7 @@ from __future__ import annotations
from dataclasses import dataclass, field
from datetime import datetime, timezone
from typing import TYPE_CHECKING, Any
from typing import TYPE_CHECKING, Any, Protocol
from lfx.schema.workflow import (
ComponentOutput,
@ -41,7 +41,18 @@ from lfx.schema.workflow import (
if TYPE_CHECKING:
from lfx.graph.graph.base import Graph
from langflow.api.v1.schemas import RunResponse
class RunResponseLike(Protocol):
"""Structural type for a v1-style run response.
Defined here so this module never imports from langflow. The backend's
``RunResponse`` (a list of ``outputs`` plus a ``session_id``) satisfies it
structurally, and lfx's own run path can supply any object with the same
shape.
"""
outputs: list[Any] | None
session_id: str | None
@dataclass(frozen=True)
@ -469,7 +480,7 @@ def _resolve_output(outputs: dict[str, ComponentOutput], selected_ids: list[str]
def run_response_to_workflow_response(
run_response: RunResponse,
run_response: RunResponseLike,
flow_id: str,
job_id: str,
inputs: dict[str, Any],

View File

@ -10,7 +10,7 @@ from __future__ import annotations
import json
import pytest
from langflow.api.v2.adapters import (
from lfx.workflow.adapters import (
StreamAdapterContext,
get_stream_adapter,
)

View File

@ -10,7 +10,7 @@ from __future__ import annotations
import json
import pytest
from langflow.api.v2.adapters import (
from lfx.workflow.adapters import (
StreamAdapterContext,
get_stream_adapter,
)

View File

@ -9,7 +9,7 @@ registers under a string name; the endpoint dispatches by name and returns
from __future__ import annotations
import pytest
from langflow.api.v2.adapters import (
from lfx.workflow.adapters import (
STREAM_ADAPTERS,
StreamAdapterContext,
StreamEvent,

View File

@ -25,7 +25,7 @@ from ag_ui.core import (
ToolCallResultEvent,
ToolCallStartEvent,
)
from langflow.api.v2.agui_translator import AGUITranslator
from lfx.workflow.agui_translator import AGUITranslator
def test_run_lifecycle_emits_started_and_finished():

View File

@ -27,7 +27,15 @@ from unittest.mock import Mock
from uuid import uuid4
import pytest
from langflow.api.v2.converters import (
from lfx.schema.workflow import (
ComponentOutput,
ErrorDetail,
JobStatus,
OutputReason,
WorkflowExecutionResponse,
WorkflowJobResponse,
)
from lfx.workflow.converters import (
_build_metadata_for_non_output,
_extract_file_path,
_extract_model_source,
@ -42,14 +50,6 @@ from langflow.api.v2.converters import (
create_job_response,
run_response_to_workflow_response,
)
from lfx.schema.workflow import (
ComponentOutput,
ErrorDetail,
JobStatus,
OutputReason,
WorkflowExecutionResponse,
WorkflowJobResponse,
)
def _setup_graph_get_vertex(graph: Mock, vertices: list[Mock]) -> None:
@ -1127,8 +1127,8 @@ class TestParseWorkflowRunRequest:
"""``parse_workflow_run_request`` projects ``WorkflowRunRequest`` onto ``ParsedWorkflowRun``."""
def test_minimal_body_round_trips_with_defaults(self):
from langflow.api.v2.converters import parse_workflow_run_request
from lfx.schema.workflow import WorkflowRunRequest
from lfx.workflow.converters import parse_workflow_run_request
parsed = parse_workflow_run_request(WorkflowRunRequest(flow_id=_VALID_UUID))
@ -1144,8 +1144,8 @@ class TestParseWorkflowRunRequest:
assert parsed.files is None
def test_full_body_round_trip(self):
from langflow.api.v2.converters import parse_workflow_run_request
from lfx.schema.workflow import WorkflowMode, WorkflowRunRequest
from lfx.workflow.converters import parse_workflow_run_request
request = WorkflowRunRequest(
flow_id=_VALID_UUID,
@ -1174,8 +1174,8 @@ class TestParseWorkflowRunRequest:
def test_run_id_is_always_none_on_the_parsed_record(self):
"""The endpoint generates run_id; callers cannot supply it via the body."""
from langflow.api.v2.converters import parse_workflow_run_request
from lfx.schema.workflow import WorkflowRunRequest
from lfx.workflow.converters import parse_workflow_run_request
parsed = parse_workflow_run_request(WorkflowRunRequest(flow_id=_VALID_UUID))
assert parsed.run_id is None