new wf_server: local static workflow server; docs
This commit is contained in:
@@ -2,6 +2,7 @@ from __future__ import annotations
|
||||
|
||||
from .listing import matches_query, paged_list_payload
|
||||
from .artifacts import WorkflowArtifactApi
|
||||
from .local_sources import builtin_sources, get_qualified_spec, qualify_spec
|
||||
from .models import RawWorkflowPlan, TraceRange
|
||||
from .capabilities import WorkflowCapabilityApi
|
||||
from .constants import (
|
||||
@@ -43,8 +44,11 @@ from .durable_context import durable_workflow_api, require_workflow_stores
|
||||
|
||||
__all__ = [
|
||||
"DEFAULT_CALL_STEP_ID",
|
||||
"builtin_sources",
|
||||
"get_qualified_spec",
|
||||
"matches_query",
|
||||
"paged_list_payload",
|
||||
"qualify_spec",
|
||||
"DEFAULT_ERROR_OUTCOME",
|
||||
"DEFAULT_ERROR_STEP_ID",
|
||||
"DEFAULT_OK_OUTCOME",
|
||||
|
||||
@@ -0,0 +1,186 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from collections.abc import Mapping
|
||||
from typing import TYPE_CHECKING, Any
|
||||
|
||||
from wf_authoring import NodeSpec, coalesce, concat, constant, default_if_none
|
||||
from wf_authoring import extract_field, filter_items, filter_items_present, first_item
|
||||
from wf_authoring import first_item_maybe, first_item_or_none, is_empty, last_item
|
||||
from wf_authoring import last_item_or_none, length, node, pick_key, pick_path
|
||||
from wf_authoring import project_fields, rename_fields, runtime_error, truthy
|
||||
from wf_authoring import extract_text_content
|
||||
from wf_core.runtime.ops.merges import DEFAULT_REDUCER_DEFINITIONS
|
||||
from wf_platform import (
|
||||
CapabilityBuckets,
|
||||
CapabilitySource,
|
||||
SourcePermissions,
|
||||
SourceVisibility,
|
||||
)
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from wf_core import ReducerSpec
|
||||
|
||||
BUILTIN_SOURCE_ID = "wf.std"
|
||||
"""Internal source id for workflow standard-library node specs."""
|
||||
|
||||
BUILTIN_CONNECTION_ID = BUILTIN_SOURCE_ID
|
||||
"""Compatibility alias for older MCP broker code."""
|
||||
|
||||
MCP_SOURCE_ID = "wf.mcp"
|
||||
"""Reserved source id for future workflow-safe MCP utility node specs."""
|
||||
|
||||
RECIPE_SOURCE_ID = "wf.recipes"
|
||||
"""Internal source id for first-party composed workflow recipes."""
|
||||
|
||||
AUTHORING_STD_SPECS: tuple[NodeSpec[Any, Any], ...] = (
|
||||
coalesce,
|
||||
default_if_none,
|
||||
constant,
|
||||
pick_key,
|
||||
pick_path,
|
||||
project_fields,
|
||||
rename_fields,
|
||||
truthy,
|
||||
runtime_error,
|
||||
first_item,
|
||||
first_item_or_none,
|
||||
first_item_maybe,
|
||||
last_item,
|
||||
last_item_or_none,
|
||||
length,
|
||||
is_empty,
|
||||
filter_items,
|
||||
filter_items_present,
|
||||
extract_field,
|
||||
concat,
|
||||
)
|
||||
"""Existing authoring ops exposed through the workflow stdlib."""
|
||||
|
||||
RECIPE_SPECS: tuple[NodeSpec[Any, Any], ...] = (extract_text_content,)
|
||||
"""Composed first-party recipes exposed as workflow-facing capabilities."""
|
||||
|
||||
|
||||
def qualify_node_name(source_id: str, local_name: str) -> str:
|
||||
"""Return one source-qualified node name without assuming MCP connections."""
|
||||
if not source_id:
|
||||
raise ValueError("source_id must not be empty")
|
||||
if not local_name:
|
||||
raise ValueError("local node name must not be empty")
|
||||
return f"{source_id}.{local_name}"
|
||||
|
||||
|
||||
def qualify_spec(source_id: str, spec: NodeSpec[Any, Any]) -> NodeSpec[Any, Any]:
|
||||
"""Return a copy of a spec with its node name scoped to a source."""
|
||||
return NodeSpec(
|
||||
name=qualify_node_name(source_id, spec.name),
|
||||
input_model=spec.input_model,
|
||||
output_model=spec.output_model,
|
||||
outcomes=spec.outcomes,
|
||||
fn=spec.fn,
|
||||
description=spec.description,
|
||||
is_async=spec.is_async,
|
||||
accepts_context=spec.accepts_context,
|
||||
input_schema_contract=spec.input_schema_contract,
|
||||
output_schema_contract=spec.output_schema_contract,
|
||||
)
|
||||
|
||||
|
||||
def get_qualified_spec(
|
||||
sources: Mapping[str, CapabilitySource],
|
||||
qualified_name: str,
|
||||
) -> NodeSpec[Any, Any]:
|
||||
"""Resolve a namespaced node spec from enabled planner-visible sources."""
|
||||
for source in sources.values():
|
||||
if not source.enabled or not source.visibility.planner:
|
||||
continue
|
||||
spec = source.capabilities.node_specs.get(qualified_name)
|
||||
if spec is not None:
|
||||
return spec
|
||||
raise KeyError(f"unknown qualified node {qualified_name!r}")
|
||||
|
||||
|
||||
def _qualified_specs(
|
||||
source_id: str,
|
||||
specs: tuple[NodeSpec[Any, Any], ...],
|
||||
) -> dict[str, NodeSpec[Any, Any]]:
|
||||
"""Return specs with authoring names rewritten under one source id."""
|
||||
local_specs = [
|
||||
node(spec, name=spec.name.removeprefix("authoring.")) for spec in specs
|
||||
]
|
||||
qualified_specs = [qualify_spec(source_id, spec) for spec in local_specs]
|
||||
return {spec.name: spec for spec in qualified_specs}
|
||||
|
||||
|
||||
def builtin_specs() -> dict[str, NodeSpec[Any, Any]]:
|
||||
"""Return primitive built-in NodeSpecs available to raw workflow plans."""
|
||||
return _qualified_specs(BUILTIN_SOURCE_ID, AUTHORING_STD_SPECS)
|
||||
|
||||
|
||||
def recipe_specs() -> dict[str, NodeSpec[Any, Any]]:
|
||||
"""Return composed first-party recipe specs."""
|
||||
return _qualified_specs(RECIPE_SOURCE_ID, RECIPE_SPECS)
|
||||
|
||||
|
||||
def builtin_reducers() -> dict[str, ReducerSpec]:
|
||||
"""Return built-in reducers owned by the workflow standard library."""
|
||||
return {
|
||||
definition.spec.name: definition.spec
|
||||
for definition in DEFAULT_REDUCER_DEFINITIONS.values()
|
||||
}
|
||||
|
||||
|
||||
def builtin_reducer_definitions():
|
||||
"""Return executable built-in reducers for trusted runtime dependency wiring."""
|
||||
return dict(DEFAULT_REDUCER_DEFINITIONS)
|
||||
|
||||
|
||||
def builtin_sources() -> dict[str, CapabilitySource]:
|
||||
"""Return all local workflow-facing capability sources."""
|
||||
return {
|
||||
BUILTIN_SOURCE_ID: CapabilitySource(
|
||||
id=BUILTIN_SOURCE_ID,
|
||||
kind="system",
|
||||
capabilities=CapabilityBuckets(
|
||||
node_specs=builtin_specs(),
|
||||
reducers=builtin_reducers(),
|
||||
reducer_definitions=builtin_reducer_definitions(),
|
||||
),
|
||||
visibility=SourceVisibility(
|
||||
planner=True,
|
||||
mcp_client=True,
|
||||
admin_dashboard=True,
|
||||
),
|
||||
permissions=SourcePermissions(safe_for_workflow=True),
|
||||
description="Workflow standard-library nodes.",
|
||||
),
|
||||
RECIPE_SOURCE_ID: CapabilitySource(
|
||||
id=RECIPE_SOURCE_ID,
|
||||
kind="system",
|
||||
capabilities=CapabilityBuckets(node_specs=recipe_specs()),
|
||||
visibility=SourceVisibility(
|
||||
planner=True,
|
||||
mcp_client=True,
|
||||
admin_dashboard=True,
|
||||
),
|
||||
permissions=SourcePermissions(safe_for_workflow=True),
|
||||
description="First-party workflow recipes composed from standard nodes.",
|
||||
),
|
||||
}
|
||||
|
||||
|
||||
__all__ = [
|
||||
"AUTHORING_STD_SPECS",
|
||||
"BUILTIN_CONNECTION_ID",
|
||||
"BUILTIN_SOURCE_ID",
|
||||
"MCP_SOURCE_ID",
|
||||
"RECIPE_SOURCE_ID",
|
||||
"RECIPE_SPECS",
|
||||
"builtin_reducer_definitions",
|
||||
"builtin_reducers",
|
||||
"builtin_sources",
|
||||
"builtin_specs",
|
||||
"get_qualified_spec",
|
||||
"qualify_node_name",
|
||||
"qualify_spec",
|
||||
"recipe_specs",
|
||||
]
|
||||
@@ -1,133 +1,41 @@
|
||||
"""Compatibility exports for workflow local sources.
|
||||
|
||||
Canonical local workflow source helpers live in `wf_api.local_sources` so
|
||||
non-MCP process hosts can construct `wf.std` without importing broker internals.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import TYPE_CHECKING, Any
|
||||
|
||||
from wf_authoring import NodeSpec, coalesce, concat, constant, default_if_none
|
||||
from wf_authoring import extract_field, filter_items, filter_items_present, first_item
|
||||
from wf_authoring import first_item_maybe, first_item_or_none, is_empty, last_item
|
||||
from wf_authoring import last_item_or_none, length, node, pick_key, pick_path
|
||||
from wf_authoring import project_fields, rename_fields, runtime_error, truthy
|
||||
from wf_authoring import extract_text_content
|
||||
from wf_core.runtime.ops.merges import DEFAULT_REDUCER_DEFINITIONS
|
||||
|
||||
from wf_platform import (
|
||||
CapabilityBuckets,
|
||||
CapabilitySource,
|
||||
SourcePermissions,
|
||||
SourceVisibility,
|
||||
from wf_api.local_sources import (
|
||||
AUTHORING_STD_SPECS,
|
||||
BUILTIN_CONNECTION_ID,
|
||||
BUILTIN_SOURCE_ID,
|
||||
MCP_SOURCE_ID,
|
||||
RECIPE_SOURCE_ID,
|
||||
RECIPE_SPECS,
|
||||
builtin_reducer_definitions,
|
||||
builtin_reducers,
|
||||
builtin_sources,
|
||||
builtin_specs,
|
||||
get_qualified_spec,
|
||||
qualify_node_name,
|
||||
qualify_spec,
|
||||
recipe_specs,
|
||||
)
|
||||
from wf_mcp.broker.service.specs import qualify_spec
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from wf_core import ReducerSpec
|
||||
|
||||
BUILTIN_CONNECTION_ID = "wf.std"
|
||||
"""Internal source id for workflow standard-library node specs."""
|
||||
|
||||
MCP_SOURCE_ID = "wf.mcp"
|
||||
"""Reserved source id for future workflow-safe MCP utility node specs."""
|
||||
|
||||
RECIPE_SOURCE_ID = "wf.recipes"
|
||||
"""Internal source id for first-party composed workflow recipes."""
|
||||
|
||||
|
||||
AUTHORING_STD_SPECS: tuple[NodeSpec[Any, Any], ...] = (
|
||||
coalesce,
|
||||
default_if_none,
|
||||
constant,
|
||||
pick_key,
|
||||
pick_path,
|
||||
project_fields,
|
||||
rename_fields,
|
||||
truthy,
|
||||
runtime_error,
|
||||
first_item,
|
||||
first_item_or_none,
|
||||
first_item_maybe,
|
||||
last_item,
|
||||
last_item_or_none,
|
||||
length,
|
||||
is_empty,
|
||||
filter_items,
|
||||
filter_items_present,
|
||||
extract_field,
|
||||
concat,
|
||||
)
|
||||
"""Existing authoring ops that are also exposed through the workflow stdlib."""
|
||||
|
||||
|
||||
RECIPE_SPECS: tuple[NodeSpec[Any, Any], ...] = (extract_text_content,)
|
||||
"""Composed first-party recipes exposed as capabilities."""
|
||||
|
||||
|
||||
def _qualified_specs(
|
||||
source_id: str,
|
||||
specs: tuple[NodeSpec[Any, Any], ...],
|
||||
) -> dict[str, NodeSpec[Any, Any]]:
|
||||
"""Return specs with authoring names rewritten under one source id."""
|
||||
local_specs = [
|
||||
node(spec, name=spec.name.removeprefix("authoring.")) for spec in specs
|
||||
]
|
||||
qualified_specs = [qualify_spec(source_id, spec) for spec in local_specs]
|
||||
return {spec.name: spec for spec in qualified_specs}
|
||||
|
||||
|
||||
def builtin_specs() -> dict[str, NodeSpec[Any, Any]]:
|
||||
"""Return primitive built-in NodeSpecs available to raw broker workflow plans."""
|
||||
return _qualified_specs(BUILTIN_CONNECTION_ID, AUTHORING_STD_SPECS)
|
||||
|
||||
|
||||
def recipe_specs() -> dict[str, NodeSpec[Any, Any]]:
|
||||
"""Return composed first-party recipe specs.
|
||||
|
||||
Recipes are wrapper-node subgraphs today. They are useful workflow-facing
|
||||
capabilities, but parent runs do not yet see their child graph frames.
|
||||
"""
|
||||
return _qualified_specs(RECIPE_SOURCE_ID, RECIPE_SPECS)
|
||||
|
||||
|
||||
def builtin_reducers() -> dict[str, ReducerSpec]:
|
||||
"""Return built-in reducers owned by the workflow standard library."""
|
||||
return {
|
||||
definition.spec.name: definition.spec
|
||||
for definition in DEFAULT_REDUCER_DEFINITIONS.values()
|
||||
}
|
||||
|
||||
|
||||
def builtin_reducer_definitions():
|
||||
"""Return executable built-in reducers for trusted runtime dependency wiring."""
|
||||
return dict(DEFAULT_REDUCER_DEFINITIONS)
|
||||
|
||||
|
||||
def builtin_sources() -> dict[str, CapabilitySource]:
|
||||
"""Return all broker-local capability sources."""
|
||||
return {
|
||||
BUILTIN_CONNECTION_ID: CapabilitySource(
|
||||
id=BUILTIN_CONNECTION_ID,
|
||||
kind="system",
|
||||
capabilities=CapabilityBuckets(
|
||||
node_specs=builtin_specs(),
|
||||
reducers=builtin_reducers(),
|
||||
reducer_definitions=builtin_reducer_definitions(),
|
||||
),
|
||||
visibility=SourceVisibility(
|
||||
planner=True,
|
||||
mcp_client=True,
|
||||
admin_dashboard=True,
|
||||
),
|
||||
permissions=SourcePermissions(safe_for_workflow=True),
|
||||
description="Workflow standard-library nodes.",
|
||||
),
|
||||
RECIPE_SOURCE_ID: CapabilitySource(
|
||||
id=RECIPE_SOURCE_ID,
|
||||
kind="system",
|
||||
capabilities=CapabilityBuckets(node_specs=recipe_specs()),
|
||||
visibility=SourceVisibility(
|
||||
planner=True,
|
||||
mcp_client=True,
|
||||
admin_dashboard=True,
|
||||
),
|
||||
permissions=SourcePermissions(safe_for_workflow=True),
|
||||
description="First-party workflow recipes composed from standard nodes.",
|
||||
),
|
||||
}
|
||||
__all__ = [
|
||||
"AUTHORING_STD_SPECS",
|
||||
"BUILTIN_CONNECTION_ID",
|
||||
"BUILTIN_SOURCE_ID",
|
||||
"MCP_SOURCE_ID",
|
||||
"RECIPE_SOURCE_ID",
|
||||
"RECIPE_SPECS",
|
||||
"builtin_reducer_definitions",
|
||||
"builtin_reducers",
|
||||
"builtin_sources",
|
||||
"builtin_specs",
|
||||
"get_qualified_spec",
|
||||
"qualify_node_name",
|
||||
"qualify_spec",
|
||||
"recipe_specs",
|
||||
]
|
||||
|
||||
@@ -0,0 +1,19 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from .context import (
|
||||
InMemoryWorkflowEventRecorder,
|
||||
LocalWorkflowRuntimeRunner,
|
||||
StaticWorkflowSpecProvider,
|
||||
WorkflowServer,
|
||||
WorkflowServerConfig,
|
||||
build_local_static_workflow_server,
|
||||
)
|
||||
|
||||
__all__ = [
|
||||
"InMemoryWorkflowEventRecorder",
|
||||
"LocalWorkflowRuntimeRunner",
|
||||
"StaticWorkflowSpecProvider",
|
||||
"WorkflowServer",
|
||||
"WorkflowServerConfig",
|
||||
"build_local_static_workflow_server",
|
||||
]
|
||||
@@ -0,0 +1,273 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from collections.abc import Mapping
|
||||
from dataclasses import dataclass, field
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
from wf_api import WorkflowApi, durable_workflow_api
|
||||
from wf_api.local_sources import builtin_sources, get_qualified_spec
|
||||
from wf_api.models import RawWorkflowPlan, TraceRange
|
||||
from wf_api.operation_context import (
|
||||
WorkflowEventRecorder,
|
||||
WorkflowOperationContext,
|
||||
WorkflowRuntimeRunner,
|
||||
WorkflowSpecProvider,
|
||||
)
|
||||
from wf_api.runtime_dependencies import resolve_runtime_dependencies
|
||||
from wf_api.saved_subgraphs import (
|
||||
SavedSubgraphTree,
|
||||
prepare_saved_subgraphs,
|
||||
resolve_saved_subgraph_tree,
|
||||
)
|
||||
from wf_api.stores import WorkflowStores, file_workflow_stores
|
||||
from wf_artifacts import WorkflowArtifact, WorkflowDeployment
|
||||
from wf_authoring import NodeSpec
|
||||
from wf_core import (
|
||||
NodeUse,
|
||||
RunState,
|
||||
Workflow,
|
||||
execute_workflow_result_async,
|
||||
resume_workflow_result_async,
|
||||
)
|
||||
from wf_platform import CapabilitySource
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class WorkflowServerConfig:
|
||||
"""Configuration for the first local/static workflow server slice."""
|
||||
|
||||
store_root: Path
|
||||
|
||||
|
||||
@dataclass(slots=True)
|
||||
class InMemoryWorkflowEventRecorder(WorkflowEventRecorder):
|
||||
"""Small process-local event sink for server composition tests."""
|
||||
|
||||
events: list[dict[str, Any]] = field(default_factory=list)
|
||||
|
||||
def record_event(self, event: object) -> None:
|
||||
self.events.append({"kind": "adapter_event", "event": event})
|
||||
|
||||
def record_workflow_event(
|
||||
self,
|
||||
event_type: str,
|
||||
*,
|
||||
capability_id: str,
|
||||
payload: dict[str, Any],
|
||||
) -> None:
|
||||
self.events.append(
|
||||
{
|
||||
"kind": event_type,
|
||||
"capability_id": capability_id,
|
||||
"payload": payload,
|
||||
}
|
||||
)
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class StaticWorkflowSpecProvider(WorkflowSpecProvider):
|
||||
"""Source provider for local/static server capabilities."""
|
||||
|
||||
sources: Mapping[str, CapabilitySource]
|
||||
|
||||
@property
|
||||
def capability_sources(self) -> dict[str, CapabilitySource]:
|
||||
return dict(self.sources)
|
||||
|
||||
def get_qualified_spec(self, qualified_name: str) -> NodeSpec[Any, Any]:
|
||||
return get_qualified_spec(self.sources, qualified_name)
|
||||
|
||||
|
||||
@dataclass(slots=True)
|
||||
class LocalWorkflowRuntimeRunner(WorkflowRuntimeRunner):
|
||||
"""Run workflow plans against local/static source catalogs."""
|
||||
|
||||
specs: StaticWorkflowSpecProvider
|
||||
artifact_store: Any
|
||||
|
||||
def compile_plan(
|
||||
self,
|
||||
plan: RawWorkflowPlan,
|
||||
node_name_bindings: dict[str, str] | None = None,
|
||||
) -> Workflow:
|
||||
node_defs: dict[str, Any] = {}
|
||||
bindings = node_name_bindings or {}
|
||||
for step in plan.nodes:
|
||||
if not isinstance(step, NodeUse):
|
||||
continue
|
||||
qualified_name = bindings.get(step.node, step.node)
|
||||
spec = self.specs.get_qualified_spec(qualified_name)
|
||||
node_defs[qualified_name] = spec.to_node_def()
|
||||
|
||||
nodes = []
|
||||
for node in plan.nodes:
|
||||
node_payload = node.model_dump(by_alias=True)
|
||||
if isinstance(node, NodeUse):
|
||||
node_payload["node"] = bindings.get(node.node, node.node)
|
||||
nodes.append(node_payload)
|
||||
|
||||
return Workflow.model_validate(
|
||||
{
|
||||
"name": plan.name,
|
||||
"input_schema": plan.input_schema,
|
||||
"state_schema": plan.state_schema,
|
||||
"output_schema": plan.output_schema,
|
||||
"output": [binding.model_dump(mode="json") for binding in plan.output],
|
||||
"outcomes": plan.outcomes,
|
||||
"start": plan.start,
|
||||
"node_defs": [node.model_dump() for node in node_defs.values()],
|
||||
"nodes": nodes,
|
||||
"edges": [edge.model_dump(by_alias=True) for edge in plan.edges],
|
||||
}
|
||||
)
|
||||
|
||||
def prepare_workflow_runtime(
|
||||
self,
|
||||
plan: RawWorkflowPlan,
|
||||
*,
|
||||
deployment: WorkflowDeployment | None,
|
||||
artifact: WorkflowArtifact | None,
|
||||
saved_subgraph_tree: SavedSubgraphTree | None = None,
|
||||
) -> tuple[Workflow, dict[str, Any], dict[str, Any], dict[str, Any]]:
|
||||
plan_node_names = [
|
||||
node.node for node in plan.nodes if isinstance(node, NodeUse)
|
||||
]
|
||||
runtime_artifact = artifact or WorkflowArtifact(
|
||||
id=plan.name,
|
||||
version=1,
|
||||
title=plan.name,
|
||||
input_schema=plan.input_schema,
|
||||
output_schema=plan.output_schema,
|
||||
outcomes=("completed",),
|
||||
plan=plan.model_dump(mode="json", by_alias=True),
|
||||
)
|
||||
dependencies = resolve_runtime_dependencies(
|
||||
artifact=runtime_artifact,
|
||||
deployment=deployment,
|
||||
sources=self.specs.capability_sources,
|
||||
plan_node_names=plan_node_names,
|
||||
)
|
||||
prepared_subgraphs = {}
|
||||
if saved_subgraph_tree is not None:
|
||||
prepared_subgraphs = prepare_saved_subgraphs(
|
||||
tree=saved_subgraph_tree,
|
||||
deployment=deployment,
|
||||
sources=self.specs.capability_sources,
|
||||
compile_plan=self.compile_plan,
|
||||
)
|
||||
elif artifact is not None and self.artifact_store is not None:
|
||||
tree = resolve_saved_subgraph_tree(
|
||||
root_artifact=artifact,
|
||||
artifact_store=self.artifact_store,
|
||||
)
|
||||
prepared_subgraphs = prepare_saved_subgraphs(
|
||||
tree=tree,
|
||||
deployment=deployment,
|
||||
sources=self.specs.capability_sources,
|
||||
compile_plan=self.compile_plan,
|
||||
)
|
||||
workflow = self.compile_plan(plan, dependencies.node_name_bindings)
|
||||
return (
|
||||
workflow,
|
||||
dependencies.node_registry,
|
||||
dependencies.reducers,
|
||||
prepared_subgraphs,
|
||||
)
|
||||
|
||||
async def run_workflow_from_plan(
|
||||
self,
|
||||
plan: RawWorkflowPlan,
|
||||
workflow_input: dict[str, Any],
|
||||
deployment: WorkflowDeployment | None = None,
|
||||
artifact: WorkflowArtifact | None = None,
|
||||
saved_subgraph_tree: SavedSubgraphTree | None = None,
|
||||
) -> RunState:
|
||||
workflow, registry, reducers, prepared_subgraphs = (
|
||||
self.prepare_workflow_runtime(
|
||||
plan,
|
||||
deployment=deployment,
|
||||
artifact=artifact,
|
||||
saved_subgraph_tree=saved_subgraph_tree,
|
||||
)
|
||||
)
|
||||
return await execute_workflow_result_async(
|
||||
workflow,
|
||||
workflow_input,
|
||||
registry,
|
||||
reducers=reducers,
|
||||
subgraphs=prepared_subgraphs,
|
||||
)
|
||||
|
||||
async def resume_workflow_from_plan(
|
||||
self,
|
||||
plan: RawWorkflowPlan,
|
||||
run: RunState,
|
||||
*,
|
||||
resume_payload: dict[str, Any],
|
||||
resume_outcome: str,
|
||||
deployment: WorkflowDeployment | None = None,
|
||||
artifact: WorkflowArtifact | None = None,
|
||||
saved_subgraph_tree: SavedSubgraphTree | None = None,
|
||||
) -> RunState:
|
||||
workflow, registry, reducers, prepared_subgraphs = (
|
||||
self.prepare_workflow_runtime(
|
||||
plan,
|
||||
deployment=deployment,
|
||||
artifact=artifact,
|
||||
saved_subgraph_tree=saved_subgraph_tree,
|
||||
)
|
||||
)
|
||||
return await resume_workflow_result_async(
|
||||
workflow,
|
||||
run,
|
||||
registry,
|
||||
resume_payload=resume_payload,
|
||||
resume_outcome=resume_outcome,
|
||||
reducers=reducers,
|
||||
subgraphs=prepared_subgraphs,
|
||||
)
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class WorkflowServer:
|
||||
"""First-slice long-lived server composition without transport concerns."""
|
||||
|
||||
config: WorkflowServerConfig
|
||||
stores: WorkflowStores
|
||||
context: WorkflowOperationContext
|
||||
api: WorkflowApi
|
||||
events: InMemoryWorkflowEventRecorder
|
||||
|
||||
@staticmethod
|
||||
def trace_range(*, start: int, limit: int) -> TraceRange:
|
||||
return TraceRange(start=start, limit=limit)
|
||||
|
||||
|
||||
def build_local_static_workflow_server(root: str | Path) -> WorkflowServer:
|
||||
"""Build a durable local/static workflow server composition."""
|
||||
config = WorkflowServerConfig(store_root=Path(root))
|
||||
stores = file_workflow_stores(config.store_root)
|
||||
events = InMemoryWorkflowEventRecorder()
|
||||
specs = StaticWorkflowSpecProvider(builtin_sources())
|
||||
runtime = LocalWorkflowRuntimeRunner(
|
||||
specs=specs,
|
||||
artifact_store=stores.artifact_store,
|
||||
)
|
||||
context = WorkflowOperationContext(
|
||||
artifact_store=stores.artifact_store,
|
||||
draft_workspace_store=stores.draft_workspace_store,
|
||||
run_store=stores.run_store,
|
||||
events=events,
|
||||
specs=specs,
|
||||
runtime=runtime,
|
||||
live_sources=None,
|
||||
)
|
||||
api = durable_workflow_api(context)
|
||||
return WorkflowServer(
|
||||
config=config,
|
||||
stores=stores,
|
||||
context=context,
|
||||
api=api,
|
||||
events=events,
|
||||
)
|
||||
Reference in New Issue
Block a user