Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
65 commits
Select commit Hold shift + click to select a range
07365da
feat(workflow): add aggregated local execution
furionw Aug 10, 2026
9310c7d
fix(workflow): enforce JSON value model
furionw Aug 12, 2026
52e912e
refactor(workflow): bind compiled plans at runtime
furionw Aug 12, 2026
c03dc69
docs(workflow): name runtime binding explicitly
furionw Aug 12, 2026
9a5d95d
feat(workflow): default compilation to local bindings
furionw Aug 12, 2026
3daab7a
feat(workflow): execute graph and handler definitions
furionw Aug 13, 2026
0da6093
fix(workflow): type unified execution branches
furionw Aug 13, 2026
89b695f
refactor(workflow): keep local execution declarative
furionw Aug 14, 2026
a84163e
refactor(workflow): align local execution authoring
furionw Aug 14, 2026
96f2e3c
refactor(workflow): name in-process orchestration explicitly
furionw Aug 15, 2026
ecf515c
fix(workflow): reject runtime stream values explicitly
furionw Aug 15, 2026
e0c816f
refactor(workflow): derive dataflow from workflow IR
furionw Aug 18, 2026
d477277
refactor(workflow): treat runtime values as opaque
furionw Aug 21, 2026
e56b286
refactor(workflow): narrow local execution lifecycle
furionw Aug 23, 2026
ebe6107
docs(workflow): mark imperative orchestration as future work
furionw Aug 23, 2026
1498696
refactor(workflow): resolve contracts in dispatcher
furionw Aug 23, 2026
8c73d97
fix(workflow): harden local execution edge cases
furionw Aug 23, 2026
b221602
refactor(workflow): preserve experimental API boundary
furionw Aug 25, 2026
0e9a409
feat(workflow): compile transport-safe remote plans
furionw Aug 14, 2026
7e0d17e
refactor(workflow): defer remote value transport checks
furionw Aug 21, 2026
acef674
style(workflow): normalize compiler imports
furionw Aug 21, 2026
a4141ea
fix(workflow): freeze authored inputs
furionw Aug 23, 2026
86c2675
fix(workflow): preserve experimental plan validation
furionw Aug 25, 2026
af1a25f
feat(workflow): execute inline stages remotely
furionw Aug 12, 2026
56769d3
style(workflow): format remote plan changes
furionw Aug 12, 2026
ac86f5d
test(workflow): keep remote proof in core fixtures
furionw Aug 12, 2026
77eaea7
feat(workflow): unify local and remote dispatch
furionw Aug 13, 2026
806d849
test(workflow): add three-process remote proof
furionw Aug 14, 2026
2197b38
refactor(workflow): simplify remote graph authoring
furionw Aug 14, 2026
0ffa716
feat(workflow): support mixed stage placement
furionw Aug 15, 2026
5b38cf6
fix(workflow): narrow remote value contracts
furionw Aug 15, 2026
4706cd2
refactor(workflow): remove remote stage envelopes
furionw Aug 20, 2026
96caf4d
refactor(workflow): make remote adapters value-opaque
furionw Aug 21, 2026
62fd2f3
feat(workflow): carry context across remote stages
furionw Aug 23, 2026
6deb756
refactor(workflow): rely on task cancellation in remote example
furionw Aug 23, 2026
5a795e9
refactor(workflow): defer remote execution example
furionw Aug 23, 2026
eb70671
refactor(workflow): derive remote context identity
furionw Aug 23, 2026
3293442
refactor(workflow): keep transport context internal
furionw Aug 23, 2026
50c0c67
style(workflow): normalize experimental imports
furionw Aug 25, 2026
a75493d
feat(workflow): serve orchestrators through model endpoints
furionw Aug 12, 2026
f221ddd
feat(workflow): bind stock Generate endpoints
furionw Aug 14, 2026
c5cac7f
test(workflow): update Generate binding authoring
furionw Aug 14, 2026
a1b62ae
refactor(workflow): separate generation streaming from collection
furionw Aug 15, 2026
c6faa16
fix(workflow): narrow generate endpoint contracts
furionw Aug 15, 2026
7d0151b
refactor(workflow): make Generate binding request-only
furionw Aug 17, 2026
ef5672a
refactor(workflow): validate Generate ports by name
furionw Aug 21, 2026
27c3767
refactor(workflow): use transport cancellation for generation
furionw Aug 23, 2026
df1ee62
fix(workflow): preserve transport context in Generate calls
furionw Aug 25, 2026
2390dd1
feat(vllm): add reusable workflow components
furionw Aug 15, 2026
c2133f3
feat(vllm): separate complete and streaming workflow contracts
furionw Aug 15, 2026
71f3c3a
refactor(vllm): keep workflow component request-only
furionw Aug 17, 2026
ae578bc
refactor(vllm): publish name-only workflow contract
furionw Aug 21, 2026
225212a
refactor(workflow): move vLLM integrations under experimental
furionw Aug 25, 2026
b47dc4d
feat(workflow): transfer remote tensors with NIXL
furionw Aug 12, 2026
a880ce2
style(workflow): format NIXL plan changes
furionw Aug 12, 2026
340ec7f
docs(workflow): describe runtime-bound carriers
furionw Aug 12, 2026
b6e1bce
style(workflow): format unified NIXL runtime
furionw Aug 13, 2026
90f3020
fix(workflow): validate NIXL dispatch types
furionw Aug 13, 2026
3a33f8b
fix(workflow): type NIXL transfer lookup
furionw Aug 13, 2026
c892dc5
fix(workflow): retain uncertain NIXL leases safely
furionw Aug 14, 2026
24d1716
test(workflow): update NIXL graph authoring
furionw Aug 14, 2026
00aec7e
test(workflow): reject mixed tensor placement
furionw Aug 15, 2026
315fc75
fix(workflow): narrow NIXL value contracts
furionw Aug 15, 2026
3d926b7
refactor(workflow): select NIXL from runtime values
furionw Aug 21, 2026
9b83542
fix(workflow): preserve experimental runtime boundaries
furionw Aug 25, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
40 changes: 40 additions & 0 deletions components/src/dynamo/experimental/workflow/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,19 +4,59 @@
"""Experimental inference workflow APIs."""

from dynamo.experimental.workflow.builder import StageHandle, Workflow
from dynamo.experimental.workflow.compiler import DeploymentSpec, compile_workflow
from dynamo.experimental.workflow.endpoint import WorkflowEndpointHandler
from dynamo.experimental.workflow.ir import StageIR, WorkflowIR
from dynamo.experimental.workflow.nixl import (
NixlLeaseRegistry,
NixlTensorCarrier,
NixlTensorFanout,
NixlTensorRef,
)
from dynamo.experimental.workflow.orchestrator import WorkflowOrchestrator
from dynamo.experimental.workflow.plan import (
ExecutionPlan,
GenerateEndpointBinding,
InlineBinding,
RemoteBinding,
)
from dynamo.experimental.workflow.remote import RemoteStageClient, RemoteStageServer
from dynamo.experimental.workflow.runtime import (
StageContext,
StageRunner,
TensorCarrier,
WorkflowExecutionError,
)
from dynamo.experimental.workflow.types import (
StageContract,
ValueRef,
WorkflowValidationError,
)

__all__ = [
"DeploymentSpec",
"StageContract",
"StageHandle",
"StageIR",
"StageContext",
"StageRunner",
"TensorCarrier",
"ValueRef",
"Workflow",
"WorkflowIR",
"WorkflowEndpointHandler",
"WorkflowExecutionError",
"WorkflowOrchestrator",
"WorkflowValidationError",
"ExecutionPlan",
"GenerateEndpointBinding",
"InlineBinding",
"NixlLeaseRegistry",
"NixlTensorCarrier",
"NixlTensorFanout",
"NixlTensorRef",
"RemoteBinding",
"RemoteStageClient",
"RemoteStageServer",
"compile_workflow",
]
2 changes: 1 addition & 1 deletion components/src/dynamo/experimental/workflow/builder.py
Original file line number Diff line number Diff line change
Expand Up @@ -127,7 +127,7 @@ def build(self) -> WorkflowIR:

return WorkflowIR(
name=self._name,
inputs=self._inputs,
inputs=frozenset(self._inputs),
stages=tuple(self._stages.values()),
outputs=self._outputs,
)
Expand Down
92 changes: 92 additions & 0 deletions components/src/dynamo/experimental/workflow/compiler.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,92 @@
# SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
# SPDX-License-Identifier: Apache-2.0

"""Lower placement-neutral workflows into physical execution plans."""

from __future__ import annotations

from dataclasses import dataclass
from types import MappingProxyType
from typing import Mapping, Optional, Union

from dynamo.experimental.workflow.builder import Workflow
from dynamo.experimental.workflow.ir import WorkflowIR
from dynamo.experimental.workflow.plan import (
Binding,
ExecutionPlan,
InlineBinding,
RemoteBinding,
)
from dynamo.experimental.workflow.types import WorkflowValidationError, validate_name


@dataclass(frozen=True)
class DeploymentSpec:
"""Stage placement requested when compiling one WorkflowIR."""

bindings: Mapping[str, Binding]

def __post_init__(self) -> None:
if not isinstance(self.bindings, Mapping):
raise WorkflowValidationError("deployment bindings must be a mapping")
bindings: dict[str, Binding] = {}
for stage_id, binding in sorted(self.bindings.items()):
validate_name(stage_id, "deployment stage id")
if not isinstance(binding, (InlineBinding, RemoteBinding)):
raise WorkflowValidationError(
f"binding for stage {stage_id!r} uses an unsupported type"
)
bindings[stage_id] = binding
object.__setattr__(self, "bindings", MappingProxyType(bindings))

@staticmethod
def inline(**runner_keys: str) -> "DeploymentSpec":
"""Build bindings to runners in the orchestrator process."""

return DeploymentSpec(
bindings={
stage_id: InlineBinding(runner_key)
for stage_id, runner_key in runner_keys.items()
}
)

@staticmethod
def remote(
*, tensor_carrier: str | None = None, **endpoint_ids: str
) -> "DeploymentSpec":
"""Build round-robin bindings to discovered Dynamo endpoints."""

return DeploymentSpec(
bindings={
stage_id: RemoteBinding(endpoint_id, tensor_carrier=tensor_carrier)
for stage_id, endpoint_id in endpoint_ids.items()
}
)


def compile_workflow(
workflow: Union[Workflow, WorkflowIR],
deployment: Optional[DeploymentSpec] = None,
) -> ExecutionPlan:
"""Compile one logical workflow, defaulting every stage to inline placement."""

workflow_ir = workflow.build() if isinstance(workflow, Workflow) else workflow
if not isinstance(workflow_ir, WorkflowIR):
raise TypeError("workflow must be a Workflow or WorkflowIR")
stage_ids = tuple(stage.id for stage in workflow_ir.stages)
if deployment is None:
deployment = DeploymentSpec.inline(
**{stage_id: stage_id for stage_id in stage_ids}
)
if not isinstance(deployment, DeploymentSpec):
raise TypeError("deployment must use DeploymentSpec")

expected = set(stage_ids)
actual = set(deployment.bindings)
if actual != expected:
raise WorkflowValidationError(
"deployment bindings differ from workflow stages; "
f"missing={sorted(expected - actual)}, extra={sorted(actual - expected)}"
)

return ExecutionPlan(workflow_ir, deployment.bindings)
Loading
Loading