Repository navigation
[Refactor] Stage Proc/Client for LLM and DiT - #5441
chickeyton wants to merge 42 commits into
Conversation
|
Codex usage limits have been reached for code reviews. Please check with the admins of this repo to increase the limits by adding credits. |
Pre-check report
Dimension summary
Verdict: 1 blocking (PR title) · 1 warning (squash). Both code issues from the first run are fixed and validated. The only remaining items are process (title + squash), not code. Fixed since first run✅ Metrics-helper divergence → resolved (
|
Signed-off-by: chickeyton <ngton2014@gmail.com>
The proc's run_stage_core param is omni_coord_address, but the manager passed it under key omni_coordinator_address, so it fell into **kwargs and was forwarded into vLLM EngineCoreProc.__init__, raising TypeError. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Signed-off-by: chickeyton <ngton2014@gmail.com>
Signed-off-by: chickeyton <ngton2014@gmail.com>
Signed-off-by: chickeyton <ngton2014@gmail.com>
The stage proc/client refactor added vllm_omni/engine/stage/* as new files, so the rebase onto main applied cleanly but silently missed main's changes to the modules those files were derived from. Carry them over: stage_replica_pool.py (from engine/stage_pool.py): - replica availability tracking (_unavailable_replicas, is_replica_available, available_replica_ids, mark_replica_unavailable, release_replica_bindings) and its use in select_replica_id / pollers - add_client(..., replica_id=...) slot restore after unregister/re-register - remove_client keeps the addr->replica_id entry - native generation-token count preferred over output-derived count - output_processor.add_request rollback via remove_request on submit failure stage_diffusion_core_proc.py (from diffusion/stage_diffusion_proc.py): - _is_executor_dead uses DiffusionExecutor.is_dead - _process_request consumes step_streaming() (step() is deprecated) Also follow two symbol moves made on main: - FinalOutputModalityType -> vllm_omni.outputs.output_metadata - MultimodalOutputProcessor -> vllm_omni.outputs.output_processor Signed-off-by: chickeyton <ngton2014@gmail.com> Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
…ffer test main added test_build_add_request_message_preserves_model_intermediate_buffer asserting isinstance(request, OmniEngineCoreRequest); this branch renamed that type to StageLLMCoreRequest. The two changes were in separate hunks so the rebase applied cleanly but left a dangling reference (NameError at runtime). Signed-off-by: chickeyton <ngton2014@gmail.com> Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
…/client These orchestrator tests exercised the pre-refactor stage APIs and were already failing on clean_stage_proc (pre-rebase); they are not rebase regressions. Align them with the refactored production surface: - Import StageReplicaPool (as StagePool) instead of the legacy engine.stage_pool StagePool, so tests drive the pool the orchestrator now uses (its poll_diffusion_output returns a list, not a single output). - FakeStageClient: rename process_engine_inputs -> process_core_inputs (orchestrator calls process_core_inputs). - FakeStageClient: rename get_output_async -> get_outputs_async and add get_outputs_nowait, matching StageCoreClientBase (the new pool polls via get_outputs_async / get_outputs_nowait). Not yet run to green end-to-end: the server deploy was interrupted before the suite could be re-verified. Signed-off-by: chickeyton <ngton2014@gmail.com> Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
The old top-level stage modules were replaced by the refactored
vllm_omni/engine/stage/ subpackage; production already imports the new
symbols. Delete the now-unused legacy modules and port their test-only
consumers to the new subpackage:
stage_engine_core_proc{,_manager}.py -> stage/stage_llm_core_proc{,_manager}.py
stage_pool.py -> stage/stage_replica_pool.py (StageReplicaPool)
stage_engine_core_client.py -> stage/stage_llm_core_client.py (StageLLMCoreClient)
diffusion/stage_diffusion_client.py -> stage/stage_diffusion_core_client.py
diffusion/stage_diffusion_proc.py -> stage/stage_diffusion_core_proc.py
Test updates:
- Re-point imports, string mock.patch targets, and the stage-pool logger
name to the new subpackage.
- Rename test_stage_engine_core_client.py -> test_stage_llm_core_client.py
and test_stage_diffusion_proc.py -> test_stage_diffusion_core_proc.py.
- Adapt to real API changes in the new clients: replica_id now defaults to
None (raises if unset), _kv_sender_host -> _core_host, and the diffusion
client's typed request/batched-output wire protocol
(add_request_async(StageDiffusionCoreRequest), get_outputs_nowait()).
stage_client.py is intentionally kept: it defines the shared StageClient/
StagePoolClient Protocols still used by the live runtime layer and is not
part of the old->new duplication.
Signed-off-by: chickeyton <ngton2014@gmail.com>
stage_client.py only ever provided static typing: the StageClient/ StagePoolClient Protocols (pure annotations) and StageClientBase (an empty `pass` base). None of it carried runtime behavior, so it can be removed outright rather than relocated: - stage_runtime.py / async_omni_engine.py: StageClient/StagePoolClient annotations -> Any (the core and inline client families share no common ancestor but object, which is why the Protocol existed). - async_omni_engine.py: drop the no-op cast(StageClient, ...); use the value directly. `cast` import retained (still used elsewhere). - inline_stage_diffusion_client.py: InlineStageDiffusionClient no longer subclasses the empty StageClientBase (its __init__ never called super()). Pure typing removal; no runtime behavior change. The unused StagePoolLLMClient/StagePoolDiffusionClient Protocols go away with the file. Signed-off-by: chickeyton <ngton2014@gmail.com>
Migrate the in-process inline diffusion client onto the same typed contract as the out-of-process StageDiffusionCoreClient so the pool drives both identically, and remove the pool's isinstance special-casing. Inline client: - Subclasses StageCoreClientBase and implements the full contract with matching signatures: add_request_async(StageDiffusionCoreRequest), batched get_outputs_nowait/get_outputs_async returning StageDiffusionCoreOutputs, shutdown(timeout=None), and _engine_dead_reason. check_health keeps its active executor probe (overriding the base template). - Requests carry sampling params as the plain-dict form (via sampling_params_to_dict); reconstructed in-process. Outputs are wrapped into StageDiffusionCoreOutput (success -> output, failure -> error), matching the out-of-process wire shape. This makes inline seed-deterministic like the out-of-process path (generator recreated from seed). Pool (stage_replica_pool): - _diffusion_add_request and poll_diffusion_output drop the isinstance(StageDiffusionCoreClient) branches; both client shapes use the typed request and batched output path. Annotations: - stage_runtime.py / async_omni_engine.py: restore StageCoreClientBase on the client-plumbing signatures (undoing the earlier Any degradation), now that every stage client is a StageCoreClientBase subclass. Tests: - test_inline_stage_diffusion_client.py: drive via StageDiffusionCoreRequest and assert on StageDiffusionCoreOutputs batches. - test_orchestrator.py: FakeStageClient routes diffusion outputs through get_outputs_nowait as a StageDiffusionCoreOutputs batch. NOTE: torch is unavailable locally, so this was verified by py_compile + static cross-reference only. The diffusion/orchestrator suites and mypy must be run in CI (or on a GPU box) before merge. Signed-off-by: chickeyton <ngton2014@gmail.com>
The debranched pool forwards a single StageDiffusionCoreRequest to diffusion clients (sampling params in plain-dict form) instead of unpacked positional args, so assert on the typed request's fields. Signed-off-by: chickeyton <ngton2014@gmail.com>
Nothing imports symbols from the stage package top-level (all 75 consumers import directly from submodules), so the lazy re-export machinery — the _EXPORTS map, __getattr__/__dir__, __all__, and the TYPE_CHECKING mirror — was dead weight. Reduce __init__ to a docstring; a trivial package init still imports nothing heavy at module scope, preserving the lightweight-submodule import property. Signed-off-by: chickeyton <ngton2014@gmail.com>
Signed-off-by: chickeyton <ngton2014@gmail.com> # Conflicts: # tests/engine/test_async_omni_engine_input.py # vllm_omni/engine/stage/stage_replica_pool.py
Signed-off-by: chickeyton <ngton2014@gmail.com> # Conflicts: # vllm_omni/core/sched/omni_ar_scheduler.py # vllm_omni/core/sched/omni_generation_scheduler.py
…ean_stage_proc_rebase Signed-off-by: chickeyton <ngton2014@gmail.com>
4c1205c to
8e6af27
Compare
… clean_stage_proc_rebase Signed-off-by: chickeyton <ngton2014@gmail.com>
…ean_stage_proc_rebase Signed-off-by: chickeyton <ngton2014@gmail.com>
… clean_stage_proc_rebase Signed-off-by: chickeyton <ngton2014@gmail.com>
|
will the change the usage of diffusion DP? for example, we currently work heavily with distributed layerwise offload, it rereuiqs work with DP replica |
… clean_stage_proc_rebase Signed-off-by: chickeyton <ngton2014@gmail.com> # Conflicts: # docs/design/architecture_overview.md # vllm_omni/diffusion/stage/inline_stage_diffusion_client.py
… clean_stage_proc_rebase Signed-off-by: chickeyton <ngton2014@gmail.com> # Conflicts: # vllm_omni/engine/orchestrator.py # vllm_omni/engine/stage/stage_replica_pool.py
There is no functional nor behavioural change, just code relocation and renaming |
|
please review the agent codes by yourself with comments first |
Omni ReviewBot triage noteAutomated triage of commit
These are automated triage suggestions only — the final decision belongs to the maintainers. |
|
|
||
| async def get_outputs_async(self) -> StageEncodeCoreOutputs: | ||
| frame = await self._output.recv() | ||
| return self._decoder.decode(frame, StageEncodeCoreOutputs) |
There was a problem hiding this comment.
Use msgspec.convert(self._decoder.decode(frame), StageEncodeCoreOutputs); decode doesn't accept a type argument.
| from vllm_omni.engine.stage_client import StageClient, StagePoolClient | ||
| from vllm_omni.engine.stage_engine_core_client import StageEngineCoreClientBase | ||
| from vllm_omni.engine.stage.stage_core_client import StageCoreClientBase | ||
| from vllm_omni.engine.stage.stage_llm_core_client import StageLLMCoreClientBase |
There was a problem hiding this comment.
Update test_stage_runtime_passes_log_stats_to_llm_replica_launch to use StageLLMCoreClientBase; the stale name raises AttributeError.
Omni ReviewBot: no human activity for 7 days@chickeyton this pull request has had no human commit, comment or review since 2026-09-16. Please confirm the current plan and next step. The author or a maintainer decides whether to change the PR state. To keep it moving, any one of these is enough: push an update, reply to the open blocker, or post the current plan and timeline. |
PLEASE FILL IN THE PR DESCRIPTION HERE.
Purpose
Refactor the following files
There will be subdirectories to contains the refactored code:
Inheritance Diagrams
flowchart BT ECP["EngineCoreProc"] SLCP["StageLLMCoreProc"] SDCP["StageDiffusionCoreProc"] CEPM["CoreEngineProcManager"] SLCPM["StageLLMCoreProcManager"] SDCPM["StageDiffusionCoreProcManager"] SLCP --> ECP SLCPM --> CEPM SLCPM -. "contains" .-o SLCP SDCPM -. "contains" .-o SDCPflowchart BT AMPC["AsyncMPClient"] DAMPC["DPLBAsyncMPClient"] SCCB["StageCoreClientBase"] SLCCB["StageLLMCoreClientBase"] SLCC["StageLLMCoreClient"] DSLCC["DPLBStageLLMCoreClient"] SDCC["StageDiffusionCoreClient"] ISDC["InlineStageDiffusionClient"] SLCC --> SCCB SLCCB --> SCCB DSLCC --> SLCCB DSLCC --> DAMPC SLCC --> AMPC SLCC --> SLCCB SDCC --> SCCB ISDC --> SCCBflowchart BT ECR["EngineCoreRequest"] ECO["EngineCoreOutput"] ECOS["EngineCoreOutputs"] SCR["StageCoreRequest"] SCO["StageCoreOutput"] SCOS["StageCoreOutputs"] SLCR["StageLLMCoreRequest"] SLCO["StageLLMCoreOutput"] SLCOS["StageLLMCoreOutputs"] SDCR["StageDiffusionCoreRequest"] SDCO["StageDiffusionCoreOutput"] SDCOS["StageDiffusionCoreOutputs"] SLCR --> ECR SLCR --> SCR SLCO --> ECO SLCO --> SCO SLCOS --> ECOS SLCOS --> SCOS SDCR --> SCR SDCO --> SCO SDCOS --> SCOSChanges
stage_core_types.py
StageCoreRequest,StageCoreOutput,StageCoreOutputsStageLLMCoreRequest,StageLLMCoreOutput,StageLLMCoreOutputsStageDiffusionCoreRequest,StageDiffusionCoreOutput,StageDiffusionCoreOutputsstage_client.py -> stage_core_client.py
StagePoolClient,StagePoolLLMClient,StagePoolDiffusionClientnot used, remove themStageClientBasetoStageCoreClientBaseStageCoreClientBasestage_engine_core_client.py -> stage_llm_core_client.py
StageEngineCoreClientBasetoStageLLMCoreClientBase(subclass fromStageCoreClientBase)StageEngineCoreClienttoStageLLMCoreClientDPLBStageEngineCoreClienttoDPLBStageLLMCoreClientStageLLMCoreClientBasesubclass fromStageCoreClientBase_default_process_engine_inputsto_default_process_core_inputsas a static function ofStageLLMCoreClientBaseEngineCoreRequestandEngineCoreOutputswithStageLLMCoreRequestandStageLLMCoreOutputs_kv_sender_hostto_core_host,_resolve_contact_hostto_resolve_core_hostengine_manager(renamed toproc_manager) andcoodinatorinmake_async_mp_client:engine/init.py
OmniEngineCoreOutputandOmniEngineCoreOutputspatch.py
OmniEngineCoreOutput,OmniEngineCoreOutputsandOmniEngineCoreRequestbyStageLLMCoreOutput,StageLLMCoreOutputsandStageLLMCoreRequeststage_engine_core_proc_manager.py -> stage_llm_core_proc_manager.py
StageEngineCoreProcManagertoStageLLMCoreProcManager__init__, argumentstart_index(the configurated dp rank) andlocal_start_indexare always set to 0 by callers, to be removedlocal_engine_counttolocal_proc_countlocal_dp_rankas not needed byStageLLMCoreProcstage_engine_core_proc.py -> stage_llm_core_proc.py
StageEngineCoreProctoStageLLMCoreProclocal_dp_rankinrun_stage_coreis not used, to be removedstage_pool.py -> stage_replica_pool.py
StagePooltoStageReplicaPoolinline_stage_diffusion_client.py
InlineStageDiffusionClientsubclass fromStageClientBasestage_diffusion_client.py -> stage_diffusion_core_client.py
StageDiffusionClienttoStageDiffusionCoreClientand subclass fromStageClientBasestage_diffusion_proc.py -> stage_diffusion_core_proc.py
StageDiffusionProctoStageDiffusionCoreProcStageDiffusionProcManagerto stage_diffusion_core_proc_manager.pystage_diffusion_core_proc_manager.py
StageDiffusionProcManagertoStageDiffusionCoreProcManagerTest Plan
vLLM Version: v0.25.0
Environment
Model
ByteDance-Seed/BAGEL-7B-MoTStageLLMCoreProc.StageDiffusionCoreProc/StageDiffusionCoreClient/StageDiffusionCoreProcManager; the KV/latent handoff from Stage 0 crosses the ZMQ boundary throughSharedMemoryConnector.Replica & parallel settings
devices"0""0"(colocated with Stage 0)DiffusionParallelConfig()defaults)max_num_seqsmax_num_batched_tokensgpu_memory_utilizationenforce_eagerenable_prefix_cachingNo tensor / sequence / data parallelism was used — this is the single-GPU,
single-replica-per-stage default. The diffusion side uses
DiffusionParallelConfig()with its default (TP=1, no SP/CFG/USP parallel).
YAML config (
vllm_omni/deploy/bagel.yaml)Launch / run command
Example (1 of 10 jobs)
Input text (edit instruction):
Input image —
bagel_02_input.png(a green fire hydrant with a yellow cap and white leaf motif on a city sidewalk):Output image —
bagel_02.png, 512×512:Test Result
requested background/style edit applied.
(
StageDiffusionCoreProc/Client/ProcManager) and Thinker stage(
StageLLMCoreProc) work end-to-end through the orchestrator after the move.relocated stage modules, and single-stage Qwen-Image-Edit = 3/3.
BEFORE SUBMITTING: read CONTRIBUTING.md and run the precheck-pr skill with the code agent for a self-check against project conventions.
(anything written below this line will be removed by GitHub Actions)