
Kdit Architecture
- 1 installs
- 62 repo stars
- Updated May 13, 2026
- tencent/ksanadit
Helps with ai & agent building tasks during AI-assisted development.
About
kdit-architecture is a Claude Code skill for ai & agent building. It helps solo builders move faster with AI-assisted coding.
- kdit-architecture
- AI & Agent Building
- AI-coding skill
Kdit Architecture by the numbers
- 1 all-time installs (skills.sh)
- Ranked #14,102 of 16,546 AI & Agent Building skills by installs in the Skillselion catalog
- Data as of Jul 20, 2026 (Skillselion catalog sync)
npx skills add https://github.com/tencent/ksanadit --skill kdit-architectureAdd your badge
Show developers this skill is listed on Skillselion. Paste this into your README.
| Installs | 1 |
|---|---|
| repo stars | ★ 62 |
| Last updated | May 13, 2026 |
| Repository | tencent/ksanadit ↗ |
What it does
Helps with ai & agent building tasks during AI-assisted development.
Files
kDiT 架构知识库
适用场景
- 理解模块关系、评估改动范围
- 架构讨论与设计评审
- 新功能涉及跨模块交互时
文档索引
| 文件 | 说明 |
|---|---|
| `overview.md` | 架构总览、组件关系、数据流、Ownership |
| `pipeline.md` | Pipeline 编排层、DAG、ExtraInputs、ContextBuilder |
| `generator.md` | Generator 去噪引擎、BaseLatent/AuxLatent、Handler |
| `node.md` | Node 计算单元、Def/Pin、dispatch_policy |
| `node-context.md` | NodeContext 上下文、metadata 禁止含 tensor |
| `pin-hub.md` | PinHub 沙箱化数据访问器 |
| `pool-key.md` | TensorPool/ModelPool、PoolKey、引用计数 |
| `adapter.md` | ComfyUI 适配层、依赖方向 |
| `device-info.md` | DeviceInfo 设备信息 |
相关 skill
| Skill | 何时调用 |
|---|---|
/kdit-standards | 编写代码时确认规范 |
/kdit-quality | 测试设计、格式检查 |
Adapter / ComfyUI
adapter 为了让 kdit 的 Node 可以适配多种其他工具的一个适配层。目前只有 ComfyUI,未来或许还有别的。
本质上,adapter 只需要一点点的适配层,就可以直接调用 kdit 的 Node 达到实现适配的目的。
语义上,ComfyUI 的 workflow 是一个 JSON 文件,包含了 Node 的组合以及 Node 之间的连接关系。这个理论上就应该可以通过Node描述成一个 PipelineDef后在本地运行。 本质 ComfyUI 的 workflow 应该和 Pipeline 在一个层级。
依赖方向: kdit/adapter/comfyui/ → kdit/ ✅; kdit/ → kdit/adapter/ ❌ 禁止
---
包结构
- `kdit/adapter/comfyui/nodes/` 是 ComfyUI 插件节点
- `kdit/nodes/` 是 kdit 内部节点
- 两者 API 完全不同,不要混淆,adapter 负责两者接口的薄适配
FeedTensor / FetchTensor 桥接
ComfyUI adapter 通过 engine.feed_tensors() 写入 tensor 到 staging area,由 FeedTensorNode(IONode)桥接到 DAG 内部。FetchTensorNode 则反向提取 DAG 输出供 ComfyUI 使用。
相关文档
- 架构总览 → `overview.md`
- Node 设计 → `node.md`
- 依赖方向与命名 → `../03_standards/naming.md`
DeviceInfo — 设备信息
源文件:`kdit/nodes/core/device_info.py`
@dataclass(frozen=True)
class DeviceInfo:
compute_device: torch.device # 计算设备
offload_device: torch.device # 卸载设备
rank_id: int # 当前卡的 rank
world_size: int # 总卡数关键规则
- DeviceInfo 嵌入
NodeContext.device字段 - 由 Executor 自动注入,禁止用户手动设置
- Pipeline 层构建 NodeContext 时
device字段为None - Executor 在调用 Node 前自动注入:
context.device = self.device_info - Node 中通过
context.device.compute_device、context.device.offload_device等访问
相关文档
- NodeContext → `node-context.md`
- 架构总览 → `overview.md`
Generator — Diffusion 去噪引擎
Generator 是 Node 内部的实现细节(被 `GeneratorNode` 封装),负责 Diffusion 去噪流程。Generator 内部的 tensor 流转(noise、denoise step 等)不受 DAG 改造影响。
---
声明式架构 (Plan E)
- `GeneratorFactory` 已废弃,新架构使用 `GeneratorDef` (frozen dataclass) + `GeneratorRunner` (final 类,无子类)
- 三个 Handler 注入模型差异: `TextHandler` / `LatentHandler` / `DenoiseHandler`
- 注册表模式: `register_generator_def()` 注册到
_GENERATOR_DEF_REGISTRY,定义文件在 `kdit/generators/defs/` 下 - `GeneratorRunner.run(ctx)` 接收 `GeneratorInferContext` 结构化上下文
---
BaseLatent 与 AuxLatent 语义规范
概述
Generator 的输入 latent 分为两类:BaseLatent(主 latent)和 AuxLatent(辅助 latent),定义在 `kdit/generators/generator_context.py`。
BaseLatent — 主 latent,决定输出尺寸
BaseLatent 是 Generator 的必需输入,其核心职责是决定 `noise_shape`,即 Generator 输入输出的主尺寸:
- 视频场景:
noise_shape决定视频的分辨率(H × W)和时长(F 帧数) - 图片场景:
noise_shape决定图片的分辨率(H × W)
noise_shape 从 base_latent.latent.shape[1:] 推导(GeneratorRunner.run() 中 noise_shape = list(base_latent_obj.latent.shape[1:]))。
# kdit/generators/generator_context.py
@dataclass
class BaseLatent:
latent: torch.Tensor # 主 tensor,noise_shape 从此推导
mask: torch.Tensor | None = None # 仅 WAN I2V / VACE 场景非 NoneBaseLatent.mask 的使用场景
| 模型 | base_latent.latent | base_latent.mask | 说明 |
|---|---|---|---|
| WAN T2V | 空 latent(torch.zeros) | None | 通过 empty latent 方法创建,仅决定 noise_shape |
| WAN I2V | VAE encode 的首帧 latent | 非 None(mask tensor) | preprocess_base_latent() 将 [mask, latent] concat |
| WAN VACE | VAE encode 的首帧 latent | 非 None(mask tensor) | 继承 WAN I2V 的行为 |
| Qwen T2I | 空 latent(torch.zeros) | None | 通过 empty latent 方法创建 |
| Qwen Edit | 空 latent(torch.zeros) | None | 通过 empty latent 方法创建 |
关键规则:目前只有 WAN I2V 及其衍生的 VACE 在传入 BaseLatent 时 mask 非 None(base_latent_list 为 [latent, mask])。其他所有模型都通过 empty latent 方法给一个空 tensor 来决定 noise_shape 和 batch_size。
AuxLatent — 辅助 latent,模型自定义用途
AuxLatent 是可选输入,理论上可以是任意 tensor 或 list,传入后由模型子类自行决定如何与 base 一起作用:
# kdit/generators/generator_context.py
@dataclass
class AuxLatent:
latent: ImageEmbeds | MultiPromptImageEmbeds | torch.Tensor类型别名:
ImageEmbeds = torch.Tensor— 单个 Tensor,shape[0]是 batch 维度MultiPromptImageEmbeds = list[torch.Tensor]— list 长度 = prompt 数量,每个 Tensor 对应一个 prompt 的参考图
各模型的 AuxLatent 用途
| 模型 | AuxLatent 内容 | 使用方式 |
|---|---|---|
| WAN T2V | None | 不使用 |
| WAN I2V | 用于噪声混合的 latent | _apply_aux_latent() 中与 noise 混合 |
| WAN VACE | 用于噪声混合的 latent | 继承 WAN I2V 的 _apply_aux_latent() |
| Qwen T2I | None | 不使用 |
| Qwen Edit | 参考图片的 VAE encode 结果 | prepare_model_forward_kargs() 中作为 ref_latents 传入模型 |
子类覆写点
| 方法 | 基类行为 | 覆写场景 |
|---|---|---|
preprocess_base_latent(base_latent_list) | 取 list[0](仅 latent) | WAN I2V: 将 [mask, latent] concat 为单个 tensor |
_apply_aux_latent(noise, aux, ...) | raise NotImplementedError | 每个子类必须实现:WAN 做噪声混合,Qwen 直接返回 noise(不混合) |
pack_aux_latent(ref_latent, patch_size) | 直接返回 | Qwen: pack_ref_latents() 做 patchify |
禁止事项
- ❌ 禁止在
BaseLatent中放非 latent 数据(如文本 embedding、配置参数) - ❌ 禁止绕过
BaseLatent直接在context.metadata中传递 noise_shape — noise_shape 必须从base_latent.latent.shape[1:]推导 - ❌ 禁止在新模型中假设
BaseLatent.mask一定为 None — 应通过preprocess_base_latent()处理 - ❌ 禁止在
AuxLatent中放不可序列化对象 — AuxLatent 通过 tensor_pool 流转,必须是 tensor 或 list[tensor]
---
相关文档
- Node 设计 → `node.md`
- 架构总览 → `overview.md`
- PinHub 沙箱 → `pin-hub.md`
- Key 类型体系 → `../03_standards/key-system.md`
NodeContext — Node 间传递的上下文
源文件:`kdit/nodes/core/node_context.py`
职责:携带配置、元数据和设备信息,可安全跨 Ray 边界序列化。
@dataclass
class NodeContext:
prompt: str | list[str] = None
negative_prompt: str | list[str] = None
img_path: str | list[str] | list[list[str]] = None
sample_config: SampleConfig = None
runtime_config: RuntimeConfig = None
cache_config: list = None
metadata: dict = field(default_factory=dict)
device: DeviceInfo = None # Executor 自动注入,禁止手动设置原则事项
- 禁止包含任何 tensor(
__post_init__校验metadata字段和一层 dict values 不含torch.Tensor) - device 字段由 Executor 自动注入,Pipeline 层构建时留
None - 可安全跨 Ray 边界序列化
- 额外配置通过
metadatadict 传递
相关文档
- 设备信息 → `device-info.md`
- 架构总览 → `overview.md`
- Node 设计 → `node.md`
Node — 计算单元
分两类:IONode(加载模型)和 InferNode(前向推理)。
---
Def/Pin 术语规范
核心概念
- Def(声明):编译时/类定义时的端口声明。Node 类上的
input_defs/output_defs。 - Pin(绑定):运行时的端口地址。Executor 构建的
input_pins/output_pins,值为 PoolKey。
命名规则
- 变量名含
def→ 声明层(TensorKey / ModelKey) - 变量名含
pin→ 运行时层(TensorPoolKey / ModelPoolKey) - Node 类属性用
input_defs/output_defs - Executor/Engine 参数用
input_pins/output_pins - PinHub 是 Pin 层的访问器
类型
| 名称 | 类型 | 层级 | 定义位置 |
|---|---|---|---|
PinDef | `TensorKey \ | ModelKey` | Def |
PinPoolKey | `TensorPoolKey \ | ModelPoolKey` | Pin |
Pins | dict[PinDef, PinPoolKey] | Pin | kdit/nodes/core/pin_def.py |
PinRef | (node_id, pin) | Pin | kdit/nodes/core/pin_def.py |
NodeRef | NodeDef 的引用 | Def | kdit/nodes/core/node_def.py |
NodeDef | frozen dataclass | Def | kdit/nodes/core/node_def.py |
---
设计原则
- Node 的入参和出参不直接传输 tensor 和 model,这两者都通过 Pool 存储,通过 PoolKey 引用。避免多卡 Ray 场景下耗时的序列化操作。
- 其他入参只能包含简单的 config 内容(通过
NodeContext),不允许存在 tensor 或 model 直接传输。 - 每个 Node 实例有唯一的
node_id(int),由模块级全局计数器(itertools.count(1))在NodeDef构造时自动分配(field(init=False)),用户不可指定。定义在kdit/nodes/core/node_def.py。计数器仅在主进程的 PipelineDefBuilder 构建 DAG 时使用,不存在多卡多进程竞争。 - DAG 中未连接的输入 pin 代表"不输入",Node 收到
None。Node 必须自行处理None输入。
---
Pin 声明
每个 Node 类通过类属性声明自己的输入输出端口(pin):
class SomeInferNode(InferNode):
input_defs = [TensorKey.POSITIVE, TensorKey.NEGATIVE] # 从 TensorPool 读的 tensor
output_defs = [TensorKey.LATENTS] # 写入 TensorPool 的 tensor- Pin 用
TensorKey/ModelKey枚举声明 IONode的 model 输出由 Factory 注册时自动填充,不需要手动声明- Pin 声明用于 build 时校验(悬空检测)和运行时 PinHub 沙箱约束
output_defs 静态声明
output_defs 保留为类属性(list[PinDef])
class GeneratorNode(InferNode):
input_defs = [TensorKey.POSITIVE, TensorKey.NEGATIVE, TensorKey.BASE_LATENT,
TensorKey.AUX_LATENT, TensorKey.VACE_CONTEXT]
output_defs = [TensorKey.LATENTS]
def run(self, pins: PinHub, *, context: NodeContext) -> None:
model = pins.get_model() # 无参 — 自动从 node_def.model_key 获取
latents = generator.run(...)
pins.put_tensor(TensorKey.LATENTS, latents)---
run() 签名
# InferNode
def run(self, pins: PinHub, *, context: NodeContext) -> None:
# IONode — 签名与 InferNode 完全一致
def run(self, pins: PinHub, *, context: NodeContext) -> None:pins是位置参数(必须),是 Node 读写数据的唯一通道context是 keyword-only,包含配置、元数据和设备信息- Node 内部禁止直接访问 TensorPool / ModelPool / DeviceInfo,全部通过
pins和context获取 - IONode 的加载参数(model_path、model_config、lora_config 等)统一放入
context.metadata
run() 返回值
| Node 类型 | 返回值 | 说明 |
|---|---|---|
IONode.run() | None | 模型写入 model_pool,不返回 |
InferNode.run() | None | 结果写入 tensor_pool,不返回 tensor |
禁止 InferNode.run() 返回 dict 或 tensor。所有中间结果通过 tensor_pool.put(key, tensor) 写入。
run() 写入规则
Node 必须无条件写入 tensor,不得判断 rank_id:
# ✅ 正确
def run(self, pins: PinHub, *, context: NodeContext) -> None:
result = unit.run(...)
pins.put_tensor(TensorKey.LATENTS, result)
# ❌ 错误 — Node 不应关心下游的 dispatch_policy
def run(self, pins: PinHub, *, context: NodeContext) -> None:
result = unit.run(...)
if context.device.rank_id == 0:
pins.put_tensor(TensorKey.LATENTS, result)---
注意事项
run()签名固定,禁止添加额外参数或**kwargs- 额外配置(包括 IONode 的加载参数)通过
context.metadata传递 - tensor 只能通过
pins.get_tensor()/pins.put_tensor()流转,禁止在参数或 metadata 中传递 tensor NodeContext.__post_init__校验 metadata 字段和一层 dict values 不含torch.Tensor
---
dispatch_policy 三维度命名
NodeDispatchPolicy 使用 input_exec_output 三维度拼接命名:
| Policy | 输入要求 | 执行范围 | 输出行为 | 典型场景 |
|---|---|---|---|---|
ALL_ALL_ALL | 所有卡都有 | 所有卡 | 各卡独立持有 | TextEncode, Generator |
R0_R0_BCAST | rank0 有即可 | 仅 rank0 | broadcast 到所有卡 | VAEEncode |
ALL_R0_R0 | 所有卡都有 | 仅 rank0 | 仅 rank0 持有 | VAEDecode |
---
Node 间数据传递
- Node 间通过 同一个 Executor 内的 tensor_pool 传递数据,不经过 Engine
- 跨 rank 数据传递由
DispatchPolicy.RANK_0_BROADCAST的 broadcast 机制自动处理 engine.get_tensor()只用于 Pipeline/ComfyUI 取回最终结果,不用于 Node 间传递- Pipeline/ComfyUI 向 Node 传递 tensor 输入时,必须通过
engine.feed_tensors()写入 tensor_pool - 禁止在 ComfyUI adapter 之间传递裸
torch.Tensor,必须传递TensorKey
---
Node 状态分析
IONode (IONode 子类)
| Node | 类变量状态 | 实例变量状态 | 说明 |
|---|---|---|---|
DiffusionLoaderNode | _pinned_memory_manager: PinnedMemoryManager (类变量) | 无 | 有类级状态:_pinned_memory_manager 是跨调用共享的单例,首次 run() 时惰性初始化 |
TextEncoderLoaderNode | 无 | 无 | 无状态 |
VAELoaderNode | 无 | 无 | 无状态 |
关键问题:DiffusionLoaderNode._pinned_memory_manager 是类变量,在 NodeFactory.create() 每次创建新实例时不会重置。这意味着:
- 同一进程内所有
DiffusionLoaderNode实例共享同一个PinnedMemoryManager - 这是有意设计(共享 pinned memory 池),但违反了"Node 无状态"的理想模型
InferNode (InferNode 子类)
| Node | 类变量状态 | 实例变量状态 | 说明 |
|---|---|---|---|
TextEncodeNode | 无 | 无 | 完全无状态 |
GeneratorNode | 无 | 无 | 完全无状态 |
VAEDecodeNode | 无 | 无 | 完全无状态 |
VAEEncodeSpatialNode | 无 | 无 | 完全无状态 |
VAEEncodeImagesNode | 无 | 无 | 完全无状态 |
所有 InferNode 完全无状态:
- 每次
run()从 PinHub 读输入 tensor 和 model - 计算结果通过
pins.put_tensor()写入 - 不持有任何跨调用的状态
---
Node 创建方式
Node 由 Factory 按 node_def 创建,Executor 内部按 node_id 缓存:
# executor.py
def _get_or_create_node(self, node_def):
if node_def.node_id in self._node_cache:
return self._node_cache[node_def.node_id]
if node_def.is_io:
node = IONodeFactory.create(node_def.node_type, node_def.model_key)
else:
node = InferNodeFactory.create(node_def.node_type, node_def.model_key)
self._node_cache[node_def.node_id] = node
return node
def run_node(self, node_def, input_pins, context):
node = self._get_or_create_node(node_def)
pin_hub = self._build_pin_hub(node_def, input_pins)
self._pre_sync_tensors(node, policy)
if is_active_rank:
node.run(pin_hub, context=context)
self._post_sync_tensors(node, node_def, policy)
self._consume_input_tensors(input_pins) # 自动消费输入 tensorNode 实例按 node_id 缓存在 Executor 中,同一 node_id 复用同一实例。
---
Node 注册
InferNode 注册
- 使用
@InferNodeFactory.register()装饰器注册 - 注册键为
(InferNodeType, [ModelKey, ...]) InferNodeType枚举值:TEXT_ENCODE,VAE_COMPUTE_SHAPE,VAE_ENCODE_SPATIAL,VAE_ENCODE_IMAGES,VAE_DECODE,GENERATE,VACE_PREPROCESS
IONode 注册
- 使用
@IONodeFactory.register()装饰器注册 - 注册键为
(IONodeType, [ModelKey | None, ...]) IONodeType枚举值:LOAD_MODEL,SAVE_VIDEO,SAVE_IMAGE,READ_IMAGE,FEED_TENSOR,FETCH_TENSOR
---
现有 Node 参考
InferNode
| Node | dispatch_policy | input_defs | output_defs |
|---|---|---|---|
T5TextEncodeNode | ALL_ALL_ALL | [] | [POSITIVE, NEGATIVE] |
QwenTextEncodeNode | ALL_ALL_ALL | [] | [POSITIVE, NEGATIVE] |
VAEEncodeSpatialNode | R0_R0_BCAST | [START_IMG, END_IMG] | [BASE_LATENT] |
VAEEncodeImagesNode | R0_R0_BCAST | [IMAGE] | [AUX_LATENT] |
VAEComputeShapeNode | R0_R0_BCAST | [] | [BASE_LATENT] |
VAEDecodeNode | ALL_R0_R0 | [LATENTS] | [VIDEO] |
GeneratorNode | ALL_ALL_ALL | [POSITIVE, NEGATIVE, BASE_LATENT, AUX_LATENT, VACE_CONTEXT] | [LATENTS] |
VACEPreprocessNode | R0_R0_BCAST | [] | [VACE_CONTEXT] |
IONode
| Node | dispatch_policy | input_defs | output_defs | 说明 |
|---|---|---|---|---|
DiffusionLoaderNode | — | — | model | 加载 Diffusion 模型 |
TextEncoderLoaderNode | — | — | model | 加载文本编码器 |
VAELoaderNode | — | — | model | 加载 VAE |
SaveVideoNode | ALL_R0_R0 | [VIDEO] | [] | 保存视频 |
SaveImageNode | ALL_R0_R0 | [VIDEO] | [] | 保存图片 |
ReadImageNode | R0_R0_BCAST | [] | [IMAGE] | 读取图片 |
FeedTensorNode | — | — | 动态 | 外部 tensor 注入桥接 |
FetchTensorNode | — | 动态 | — | 外部 tensor 提取桥接 |
---
相关文档
- 编码实操规范(InferNode 开发 checklist) → `03_standards/node-and-tensor.md`
- PinHub 沙箱机制 → `pin-hub.md`
- NodeContext 详情 → `node-context.md`
- TensorPool / ModelPool → `pool-key.md`
- Key 类型体系 → `03_standards/key-system.md`
kDiT 架构总览
组件关系图
%%{init: {'theme': 'dark'}}%%
graph TB
subgraph External ["外部入口"]
ComfyUI["ComfyUI Adapter\nkdit/adapter/comfyui/"]
UserAPI["Pipeline.generate() API"]
end
subgraph PipelineLayer ["Pipeline 编排层"]
PipelineDef["PipelineDef\nfrozen dataclass — DAG 定义\nnodes + edges"]
PipelineDefBuilder["PipelineDefBuilder\n链式构建 DAG\nadd_loader / add_infer / connect"]
Pipeline["Pipeline\nfrom_models / load_models / generate"]
ContextBuilder["ContextBuilder\n为每个 NodeDef 构建 NodeContext"]
ExtraInputs["ExtraInputs\n模型特有输入(子类化)"]
DAG["topo_sort + compute_input_pins\n拓扑排序 → 计算 input_pins"]
end
subgraph EngineLayer ["Engine 分发层"]
Engine["Engine\nClassVar 单例 (get_default)\n纯分发,不持有资源"]
AutoDispatch["@auto_dispatch\n透明单卡/多卡切换"]
end
subgraph ExecutorLayer ["Executor 执行层"]
Executor["Executor\n持有 TensorPool + ModelPool\n+ DeviceInfo + node_cache"]
RayExecutor["RayExecutor\nRay remote Actor\n继承 Executor"]
TorchDist["torchrun 模式\n单 Executor + DDP"]
end
subgraph NodeLayer ["Node 计算层"]
IONode["IONode\n模型加载\nIONodeType 枚举"]
InferNode["InferNode\n推理计算\nInferNodeType 枚举"]
PinHub["PinHub\n沙箱化数据访问器\nget_model / get_tensor / put_tensor"]
NodeContext["NodeContext\n可序列化上下文\nmetadata 禁止含 Tensor"]
DeviceInfo["DeviceInfo\nfrozen dataclass\nExecutor 注入"]
NodeDef["NodeDef\nfrozen dataclass\nnode_id(auto) + node_type + model_key\nkdit/nodes/core/node_def.py"]
DispatchPolicy["NodeDispatchPolicy\nALL_ALL_ALL\nR0_R0_BCAST\nALL_R0_R0"]
end
subgraph GeneratorLayer ["Generator 去噪引擎"]
GeneratorDef["GeneratorDef\nfrozen dataclass\nmodel_key + 3 Handlers"]
GeneratorRunner["GeneratorRunner\nfinal 类,无子类\n统一去噪主流程"]
TextHandler["TextHandler\n文本 conditioning"]
LatentHandler["LatentHandler\nlatent 预处理 / pack / unpack"]
DenoiseHandler["DenoiseHandler\n去噪循环钩子"]
GenContext["GeneratorInferContext\n模型 + tensor + 设备 + 配置"]
BaseLatent["BaseLatent\n主 latent → noise_shape"]
AuxLatent["AuxLatent\n辅助 latent(可选)"]
end
subgraph PoolLayer ["Pool 数据层"]
TensorPool["TensorPool\n存储 + 引用计数 + 自动释放"]
ModelPool["ModelPool\n模型实例管理"]
TensorPoolKey["TensorPoolKey\n(node_id, TensorKey)"]
ModelPoolKey["ModelPoolKey\n(node_id, ModelKey)"]
TensorValue["TensorValue\n包装 Tensor | list[Tensor]"]
end
subgraph KeySystem ["Key 体系"]
ModelKey["ModelKey 枚举\n模型身份标识"]
TensorKey["TensorKey 枚举\n语义 tensor 标识"]
PipelineKey["PipelineKey 枚举\nPipeline 身份标识"]
InferNodeTypeK["InferNodeType 枚举\n推理节点类型"]
IONodeTypeK["IONodeType 枚举\n加载节点类型"]
PinDef["PinDef = TensorKey | ModelKey\nPin 声明类型\nkdit/nodes/core/pin_def.py"]
end
%% 外部入口 → Pipeline
ComfyUI -->|"调用"| Pipeline
UserAPI -->|"调用"| Pipeline
%% Pipeline 内部
PipelineDefBuilder -->|".build()"| PipelineDef
Pipeline -->|"持有"| PipelineDef
Pipeline -->|"使用"| ContextBuilder
Pipeline -->|"调用"| DAG
ContextBuilder -->|"读取"| ExtraInputs
%% Pipeline → Engine
Pipeline -->|"run_node()"| Engine
%% Engine 分发
Engine -->|"单卡"| Executor
Engine -->|"Ray 多卡"| RayExecutor
Engine -->|"torchrun 多卡"| TorchDist
AutoDispatch -.->|"装饰"| Engine
%% Executor → Node
Executor -->|"_run_io_node()"| IONode
Executor -->|"_run_infer_node()"| InferNode
Executor -->|"构建"| PinHub
Executor -->|"注入"| DeviceInfo
Executor -->|"持有"| TensorPool
Executor -->|"持有"| ModelPool
%% Node 运行
IONode -->|"run(pins, context)"| PinHub
InferNode -->|"run(pins, context)"| PinHub
InferNode -->|"接收"| NodeContext
PinHub -->|"读写"| TensorPool
PinHub -->|"读写"| ModelPool
%% Generator 子系统(GeneratorNode 内部调用)
InferNode -->|"GeneratorNode 调用"| GeneratorRunner
GeneratorRunner -->|"持有"| GeneratorDef
GeneratorDef -->|"组合"| TextHandler
GeneratorDef -->|"组合"| LatentHandler
GeneratorDef -->|"组合"| DenoiseHandler
GeneratorRunner -->|"接收"| GenContext
GenContext -->|"包含"| BaseLatent
GenContext -->|"包含"| AuxLatent
%% Pool 内部
TensorPool -->|"存储"| TensorValue
TensorPool -->|"索引"| TensorPoolKey
ModelPool -->|"索引"| ModelPoolKey
%% Key 关联
TensorPoolKey -->|"组合"| TensorKey
ModelPoolKey -->|"组合"| ModelKey---
数据流
Pipeline.generate()
→ DAG topo_sort → 按拓扑序遍历 NodeDef
→ ContextBuilder.build_context(node_def, inputs) → NodeContext
→ compute_input_pins(node_def, edges, all_outputs) → input_pins
→ Engine.run_node(node_def, input_pins, context)
→ Executor.run_node(node_def, input_pins, context)
→ _get_or_create_node(node_def) → IONode | InferNode
→ _build_pin_hub(node_def, input_pins) → PinHub
→ _inject_context_defaults(node_def, context) → DeviceInfo 注入
→ _pre_sync_tensors(node, policy)
→ node.run(pins, context=context)
→ pins.get_tensor() / pins.put_tensor() ← TensorPool
→ pins.get_model() / pins.put_model() ← ModelPool
→ _post_sync_tensors(node, node_def, policy)
→ _build_output_pins(node, node_def) → output_pins
→ _consume_input_tensors(input_pins) ← 引用计数递减
← output_pins
→ all_outputs[node_def.node_id] = output_pins
→ engine.get_tensor(TensorKey.VIDEO) → 最终输出---
Ownership 与状态关系
整体 Ownership 层级
Engine (singleton via get_default / 或多实例)
├── owns: executors
│ ├── 单卡模式: 1 个 Executor 实例
│ └── 多卡模式: N 个 RayExecutor (Ray Actor)
├── owns: num_gpus, _is_ray, _cleaned_up (引擎级元数据)
├── NOT own: model_pool, tensor_pool, device 信息 (这些属于 Executor)
└── NOT own: 任何 Node 实例 (Node 由 IONodeFactory/InferNodeFactory 按需创建,缓存在 Executor 中)
Executor (每卡一个实例)
├── owns: model_pool — ModelPool (存储已加载的模型)
├── owns: tensor_pool — TensorPool (存储推理中间 tensor)
├── owns: dist_group — DistributedGroupManager (管理 torch.distributed)
├── owns: device_ctx — DeviceInfo (frozen dataclass, 只读)
├── owns: device / offload_device / device_id (设备信息)
├── owns: rank_id / world_size (分布式信息)
├── owns: dist_config / shard_fn (分布式配置)
└── owns: _node_cache — Node 实例按 node_id 缓存复用Engine (`kdit/engine/engine.py`)
| 属性 | 类型 | 说明 |
|---|---|---|
executors | Executor 或 list[RayExecutor] | 唯一核心持有物。单卡时是一个实例,多卡时是 Ray Actor 列表 |
num_gpus | int | GPU 数量,从 dist_config 复制 |
_is_ray | bool | 是否使用 Ray 分布式 |
_cleaned_up | bool | 清理标记,防止重复清理 |
Engine 不持有:model_pool、tensor_pool、device 信息、Node 实例。Engine 是纯粹的分发层,所有实际资源都在 Executor 上。
Engine 公开 API(桥接方法,透传到 Executor)
| 方法 | 用途 | 说明 |
|---|---|---|
engine.run_node() | 执行 Node | 分发到所有 Executor,返回 `output_pins`。Ray 模式取 rank 0 结果 |
engine.get_tensor(key) | 取回 TensorValue | 自动从 rank 0 取,返回 TensorValue(需 .data 取裸 tensor) |
engine.feed_tensors(tensors) | 写入 tensor | 写入所有 Executor 的 tensor_pool,自动包装为 TensorValue,返回 feed_pins |
engine.has_tensor(key) | 检查 key 存在性 | 检查 rank 0 的 tensor_pool 中是否存在指定 key |
engine.register_tensor(pool_key, ref_count) | 注册引用计数 | 透传到所有 Executor 的 tensor_pool.register() |
engine.clear_all_tensors() | 清理所有 tensor | 清理所有 Executor 的 tensor_pool — 用于 try/finally 异常恢复 |
engine.rename_tensor(old, new) | 重命名 key | 透传到所有 Executor 的 tensor_pool.rename() |
`tensor_scope` 已删除。异常安全通过 try/finally + engine.clear_all_tensors() 实现。TensorPool 内置引用计数(register/consume)自动管理中间 tensor 的释放。
Executor (`kdit/executor/executor.py`)
| 属性 | 类型 | 生命周期 | 说明 |
|---|---|---|---|
model_pool | ModelPool | 与 Executor 同生命周期 | 存储所有已加载模型,按 ModelKey 索引 |
tensor_pool | TensorPool | 每次推理结束时 clear(Pipeline 用 try/finally) | 存储推理中间 tensor,内置引用计数 |
dist_group | DistributedGroupManager | 与 Executor 同生命周期 | 管理 broadcast 等分布式操作 |
device_ctx | DeviceInfo | 初始化后不变(frozen) | 只读设备上下文,传入 Node.run() |
device | torch.device | 不变 | 计算设备 (如 cuda:0) |
offload_device | torch.device | 不变 | 卸载设备 (如 cpu) |
dist_config | DistributedConfig | init_torch_dist_group() 后更新 | 分布式配置 |
shard_fn | partial 或 None | init_torch_dist_group() 后设置 | FSDP 分片函数 |
Executor 同步机制
Executor.run_node() 负责:
1. `_pre_sync_tensors()`: 执行前的 tensor 同步(预留接口,未来可自动 broadcast 输入) 2. `is_active_rank`: 根据 policy 判断当前卡是否执行 run() 3. `_post_sync_tensors()`: 执行后的 tensor 同步(R0_R0_BCAST 时 broadcast output_defs 中的 key) 4. `_consume_input_tensors()`: 自动消费输入 tensor 引用计数
Node 内部不需要感知多卡逻辑,Executor 负责所有 tensor 的 pre/post 同步和引用计数管理。
---
关键设计决策
1. Engine 是纯分发层:不持有 model_pool / tensor_pool / device 信息,所有资源在 Executor 上 2. 多卡模式:优先检测 torchrun 环境,否则使用 Ray。Engine 透明切换 3. Generator 是 InferNode 的内部子系统:GeneratorNode.run() 内部构建 GeneratorInferContext,调用 GeneratorRunner.run(ctx) 执行去噪 4. Adapter 依赖方向:kdit/adapter/comfyui/ → kdit/ ✅;反向 ❌ 禁止 5. Node 通过 PinHub 沙箱访问数据:禁止直接操作 TensorPool / ModelPool 6. NodeContext 可序列化:metadata 禁止含 torch.Tensor,保证跨 Ray 边界安全
设计约束总结
| 约束 | 说明 |
|---|---|
| Engine 不持有资源 | Engine 只是分发层,所有实际资源(model_pool, tensor_pool, device)在 Executor 上 |
| Executor 持有所有资源 | model_pool + tensor_pool + dist_group + device_ctx + node_cache |
| Node 无状态(理想) | InferNode 完全无状态;IONode 中 DiffusionLoaderNode 有类级 _pinned_memory_manager 例外 |
| DeviceInfo 只读 | frozen=True dataclass,Node 无法篡改 |
| NodeContext 无 tensor | __post_init__ 强制校验不含 torch.Tensor,保证可跨 Ray 序列化 |
| tensor_pool 生命周期 | Pipeline 用 try/finally + clear_all_tensors();DAG 模式下引用计数自动释放 |
| model_pool 生命周期 | 与 Executor 同生命周期,需手动 clear_models() 释放 |
| InferNode.run() 签名固定 | (self, pins: PinHub, *, context: NodeContext) -> None,禁止扩展 |
| Tensor 只能通过 PinHub 流转 | 禁止在 run() 参数或 context.metadata 中传递 tensor |
---
子系统文档
| 子系统 | 文档 |
|---|---|
| Pipeline 编排层 | `pipeline.md` |
| Generator 去噪引擎 | `generator.md` |
| Node 计算单元 | `node.md` |
| PinHub 沙箱 | `pin-hub.md` |
| NodeContext | `node-context.md` |
| Pool / PoolKey | `pool-key.md` |
| Key 类型体系 | `../03_standards/key-system.md` |
| Adapter 规范 | `adapter.md` |
| 设备信息 | `device-info.md` |
| 编码实操规范 | `../03_standards/node-and-tensor.md` |
PinHub — 沙箱化数据访问器
职责:Node 运行时的数据读写路由层。每个 Node 实例拥有独立的 PinHub,被严格约束在 DAG 声明的范围内。
---
核心机制
- 读操作:根据
input_pins(DAG 连线计算的映射),将上游 Node 的输出 PoolKey 映射到当前 Node 的 input pin - 写操作:自动用
当前 node_id + pin生成 PoolKey 写入 Pool
pins.get_tensor(TensorKey.POSITIVE) # 从 input_pins 查找上游的 TensorPoolKey,读取
pins.put_tensor(TensorKey.LATENTS, data) # 写入 TensorPoolKey(self.node_id, TensorKey.LATENTS)
pins.get_model(ModelKey.T5TextEncoder) # 从 input_pins 查找上游的 ModelPoolKey,读取
pins.put_model(ModelKey.VAE_WAN2_1, model) # 写入 ModelPoolKey(self.node_id, ModelKey.VAE_WAN2_1)---
API 详细说明
Tensor 操作
| 方法 | 签名 | 说明 |
|---|---|---|
get_tensor(key) | (TensorKey) → TensorValue | 从 input_pins 查找上游 TensorPoolKey,读取 TensorValue。未连线返回 None |
peek_tensor(key) | `(TensorKey) → TensorValue \ | None` |
put_tensor(key, tensor) | (TensorKey, Tensor) → None | 写入 TensorPoolKey(self.node_id, key) 到 tensor_pool |
Model 操作
| 方法 | 签名 | 说明 |
|---|---|---|
get_model() | () → ModelBase | 无参,自动从 node_def.model_key 获取对应的 ModelPoolKey,读取模型 |
get_model(key) | (ModelKey) → ModelBase | 指定 ModelKey 读取(多模型 Node 场景) |
put_model(key, model) | (ModelKey, ModelBase) → None | 写入 ModelPoolKey(self.node_id, key) 到 model_pool |
---
沙箱约束
- 读:只能读
input_pins中存在的 key(即 DAG 连线声明的上游输出),读不到其他 Node 的数据 - 写:只能写
PoolKey(self.node_id, pin),即自己 node_id 命名空间下的 key,写不到别人的命名空间 - 未连线的 optional tensor pin 返回
None,未连线的 required model pin 抛出KeyError
---
构建位置
PinHub 在 Executor 内部构建(不在 Pipeline 层),因为:
tensor_pool和model_pool活在 Executor 上(每卡一份)- 多卡下每个 Executor 各自构建自己的 PinHub,天然正确
- PinHub 不需要跨进程序列化
---
注意事项
input_pins是纯数据(dict[TensorKey, TensorPoolKey]+dict[ModelKey, ModelPoolKey]),由 Pipeline 层从 DAG edges 计算,通过 Engine 分发到 Executor- PinHub 不持有
DeviceInfo,设备信息通过context.device获取
---
相关文档
- Node 设计原则与 Pin 声明 → `node.md`
- PoolKey 体系 → `pool-key.md`
- NodeContext → `node-context.md`
Pipeline — 编排层
职责:DAG 遍历 + 计算 input_pins + 分发执行。
PipelineDef是不可变的 DAG 定义(frozen dataclass),包含nodes+edgesPipelineDefBuilder链式构建,通过.add_loader()/.add_infer()/.connect()声明- Pipeline 层负责 DAG 拓扑排序、条件检查、构建 NodeContext
- Pipeline 层计算
input_pins(纯数据),传给 Engine → Executor - Executor 只执行单个 Node,不感知 DAG
Pipeline.generate()接收extra_inputs: ExtraInputs | None传递模型特有输入,禁止**kwargsContextBuilder是 Pipeline 和 Node 之间的桥梁,负责prepare_generate_inputs()+build_context()
NodeRef / PinRef — DAG 连线引用
NodeRef:add_loader()/add_infer()返回的 Node 引用,支持属性访问生成 PinRef。定义在kdit/nodes/core/node_def.pyPinRef:(node_id, pin)frozen dataclass,用于connect()声明连线。定义在kdit/nodes/core/pin_def.py- 旧的
kdit/pipelines/pin_ref.py已删除,NodeRef 和 PinRef 分别搬到上述位置
vae_a = builder.add_infer(InferNodeType.VAE_ENCODE_SPATIAL, ModelKey.VAE_WAN2_1)
gen = builder.add_infer(InferNodeType.GENERATE, ModelKey.Wan2_2_I2V_14B)
# vae_a.BASE_LATENT → PinRef(node_id, TensorKey.BASE_LATENT)
builder.connect((vae_a.BASE_LATENT, gen.BASE_LATENT))NodeRef.__getattr__在 TensorKey 和 ModelKey 枚举中查找属性名,不需要引号- 只有相同类型(都是 TensorKey 或都是 ModelKey)的 pin 才能 connect,但不要求同名
---
Pipeline.generate() 接口规范
签名
def generate(
self,
prompt: str | list[str],
*,
prompt_negative: str | list[str] | None = None,
sample_config: SampleConfig = None,
runtime_config: RuntimeConfig = None,
cache_config: list[CacheConfig | HybridCacheConfig] | None = None,
extra_inputs: ExtraInputs | None = None,
):规则
- 禁止
**kwargs—generate()签名中不允许出现**kwargs - 禁止在
generate()签名中添加模型特有的具名参数(如start_img_path=) - 模型特有输入必须通过
extra_inputs参数传入 - 无特有输入的模型(如 T2V、T2I)不需要传
extra_inputs
公共参数 vs 模型特有参数
| 类别 | 参数 | 说明 |
|---|---|---|
| 公共 | prompt, prompt_negative | 所有模型都需要 |
| 公共 | sample_config, runtime_config, cache_config | 采样/运行时/缓存配置 |
| 模型特有 | extra_inputs | 通过 ExtraInputs 子类传入 |
---
ExtraInputs — 模型特有输入管理
基类
# kdit/pipelines/extra_inputs.py
from dataclasses import dataclass
@dataclass
class ExtraInputs:
"""模型特有输入的基类。
每个 Pipeline 定义自己的子类。
T2V/T2I 等无特有输入的 Pipeline 不需要传此参数。
"""
pass子类定义位置
每个 ExtraInputs 子类定义在对应的 ContextBuilder 文件中:
| 子类 | 文件 | 字段 |
|---|---|---|
WanI2VExtraInputs | kdit/pipelines/context_builders/wan.py | start_img_path, end_img_path, video_control_config, aux_latent(I2V 和 VACE 共用) |
QwenEditExtraInputs | kdit/pipelines/context_builders/qwen.py | img_path |
设计原则
1. 所有字段都有默认值(通常为 None),使 ExtraInputs() 空构造合法 2. 字段类型明确,IDE 可自动补全 3. 子类只包含用户传入的原始输入(如 start_img_path: str),不包含内部中间数据(如 start_img_tensor: Tensor) 4. 内部中间数据由 ContextBuilder 的 _extra 属性管理,与用户侧 ExtraInputs 分离
用户侧调用示例
# WanI2V — 有模型特有输入
from kdit.pipelines.context_builders.wan import WanI2VExtraInputs
pipeline.generate(
prompts,
extra_inputs=WanI2VExtraInputs(
start_img_path="path/to/img.jpg",
end_img_path="path/to/end.jpg",
),
sample_config=SampleConfig(steps=40),
runtime_config=RuntimeConfig(seed=1234, size=(1280, 720), frame_num=81),
)
# WanT2V — 无特有输入,不传 extra_inputs
pipeline.generate(
prompts,
sample_config=SampleConfig(steps=40),
runtime_config=RuntimeConfig(seed=1234, size=(1280, 720), frame_num=81),
)---
ContextBuilder 开发规范
职责
ContextBuilder 是 Pipeline 和 Node 之间的桥梁:
1. `prepare_generate_inputs()` — 从 ExtraInputs 提取、校验、预处理模型特有输入 2. `build_context()` — 为每个 NodeDef 构建 NodeContext 3. `check_condition()` — 判断条件节点是否执行
prepare_generate_inputs() 签名
def prepare_generate_inputs(
self,
base_inputs: PipelineGenerateInputs,
extra_inputs: ExtraInputs | None,
*,
default_settings: Any,
engine: Engine,
vae_model_key: ModelKey | None,
) -> None:extra_inputs是用户传入的结构化输入default_settings、engine、vae_model_key是内部注入的显式参数(不再通过 kwargs)- 子类应在此方法中校验
extra_inputs类型,并将处理后的中间数据存入self._extra
build_context() 签名
@abstractmethod
def build_context(
self,
node_def: NodeDef,
inputs: PipelineGenerateInputs,
) -> NodeContext:- 参数是
NodeDef(不是InferTask) - 通过
node_def.node_type分支构建不同 Node 的 context - 通过
node_def.node_id区分同类型 Node 的不同实例
多实例 Node 的 context 区分
当 DAG 中有多个同类型 Node 实例(如两个 ReadImageNode)时,ContextBuilder 通过 DAG edges 查找 node_id 的输出连接到下游的哪个 pin,从而决定传入哪个输入:
def _build_read_image_ctx(self, node_def: NodeDef, inputs) -> NodeContext:
# 查找此 ReadImage 实例的 IMAGE 输出连接到下游的哪个 pin
dst_pin = self._find_downstream_pin(node_def.node_id, TensorKey.IMAGE)
if dst_pin == TensorKey.START_IMG:
img_paths = self._extra.start_img_path
elif dst_pin == TensorKey.END_IMG:
img_paths = self._extra.end_img_path
return NodeContext(metadata={"img_paths": img_paths})前提:ContextBuilder 初始化时需要接收 pipeline_def 以访问 edges 信息。
禁止事项
- 禁止
prepare_tensors()— 所有 tensor 注入通过 DAG Node 完成 - 禁止在
prepare_generate_inputs()中使用**kwargs— 所有参数显式声明 - 禁止在
build_context()中直接操作 tensor_pool — tensor 只能通过 Node 的 PinHub 流转
---
悬空 Pin 规则
规则
DAG 中未连接的输入 pin 代表"不输入":
- Node 通过
pins.get_tensor(key)读取时返回None - Node 必须自行处理
None输入 - 不需要自动补全、不需要 INPUT_NODE 概念
示例
# GeneratorNode 已正确处理 AUX_LATENT = None
class GeneratorNode(InferNode):
input_defs = [TensorKey.POSITIVE, TensorKey.NEGATIVE,
TensorKey.BASE_LATENT, TensorKey.AUX_LATENT,
TensorKey.VACE_CONTEXT]
def run(self, pins, *, context):
aux_latent_val = pins.get_tensor(TensorKey.AUX_LATENT)
aux_latent = aux_latent_val.data if aux_latent_val is not None else None
# aux_latent 可能为 None,后续逻辑正确处理_validate_dag() 校验
- 悬空的输入 tensor pin → 不报错,Node 收到 None
- 悬空的输入 model pin → 报错(model 是必需的)
- 重复的
(dst_id, dst_pin)→ 报错(一个输入 pin 只能有一条入边)
---
Cross-Pin Connect
规则
DAG 连线允许不同名的 pin 之间连接,只要它们是同一类型(都是 TensorKey 或都是 ModelKey):
# 合法 — IMAGE → START_IMG(都是 TensorKey)
read_s.IMAGE >> venc.START_IMG
# 合法 — IMAGE → END_IMG(都是 TensorKey)
read_e.IMAGE >> venc.END_IMG
# 非法 — TensorKey → ModelKey(类型不同)
read_s.IMAGE >> venc.MODEL # TypeError用途
Cross-pin connect 使单一功能的 Node(如 ReadImageNode 只输出 IMAGE)可以连接到不同语义的下游 pin,通过 DAG 多实例实现复用。
PoolKey — 实例化的存储 key
---
ModelPoolKey / TensorPoolKey
- `ModelPoolKey` 定义在
kdit/models/model_pool_key.py - `TensorPoolKey` 定义在
kdit/tensor/tensor_pool_key.py
# kdit/models/model_pool_key.py
@dataclass(frozen=True)
class ModelPoolKey:
node_id: int # NodeDef 全局自增自动分配的唯一 ID
pin: ModelKey # 枚举直接存储
# kdit/tensor/tensor_pool_key.py
@dataclass(frozen=True)
class TensorPoolKey:
node_id: int # NodeDef 全局自增自动分配的唯一 ID
pin: TensorKey # 枚举直接存储frozen=True自动生成__hash__和__eq__,可直接用作 dict keypin直接用枚举(不用.value),类型安全且可读- 分离 ModelPoolKey / TensorPoolKey 是因为 Model 和 Tensor 生命周期不同
- Pool 中直接传
ModelKey/TensorKey作为 key 已标记为 deprecated,新代码应使用ModelPoolKey/TensorPoolKey
为什么需要 PoolKey
解决"同一个 Pipeline 中不能有两个相同类型的 Node 实例"的问题。例如两个 VAE_ENCODE_SPATIAL 节点,各自写入 TensorPoolKey(node_id=4, TensorKey.BASE_LATENT) 和 TensorPoolKey(node_id=5, TensorKey.BASE_LATENT),互不冲突。
---
TensorPool (`kdit/tensor/tensor_pool.py`)
- Owner: Executor
- 生命周期: Pipeline 推理结束时由
engine.clear_all_tensors()清理;DAG 模式下中间 tensor 通过引用计数自动释放 - 内容:
dict[TensorPoolKey, TensorValue],每个 TensorValue 持有Tensor | list[Tensor] - 用途: Node 间通过
TensorKey引用 tensor,避免 tensor 跨 Ray 边界序列化
关键方法
| 方法 | 说明 |
|---|---|
put(key, tensor) | 写入 tensor,自动包装为 TensorValue |
get(key) | 读取 TensorValue(消费引用计数) |
peek(key) | 读取 TensorValue(不消费引用计数) |
has(key) | 检查 key 是否存在 |
clear(exclude) | 释放除 exclude 列表外的所有 tensor,重置引用计数 |
register(pool_key, ref_count) | 注册 tensor 的下游消费者数 |
consume(pool_key) | 消费一次引用计数,降为 0 时自动 release |
remove(pool_key) | 强制移除 tensor |
rename(old, new) | 重命名 key |
TensorValue 类
TensorValue 是 tensor_pool 中值的包装类,持有单个 torch.Tensor 或 list[torch.Tensor],负责释放:
class TensorValue:
__slots__ = ['data']
def __init__(self, data: torch.Tensor | list[torch.Tensor]):
self.data = data
def release(self):
"""释放持有的 tensor 引用。"""
if isinstance(self.data, list):
for i in range(len(self.data)):
self.data[i] = None
self.data.clear()
self.data = None- Pool 的
get(key)返回TensorValue,Node 内部用.data取裸 tensor - Pool 的
clear()遍历调用TensorValue.release() - 只有最终边界(如 vae_decode 输出给用户)才允许
.data取裸 tensor
引用计数机制
TensorPool 内置 register() / consume() / remove() 引用计数机制,替代了旧的 tensor_scope:
# Pipeline DAG 模式 — Executor 自动管理
# 1. Pipeline 构建 DAG 时,Engine 调用 register_tensor() 注册每个 tensor 的下游消费者数
# 2. Executor.run_node() 执行后自动 consume 输入 tensor
# 3. consume 时 ref_count 降为 0 → 自动 release TensorValue
# ComfyUI 模式 — try/finally 手动清理
try:
engine.put_tensors({TensorKey.IMAGE: image})
engine.run_node(node_def, input_pins, context)
result = engine.get_tensor(TensorKey.AUX_LATENT)
finally:
engine.clear_all_tensors()clear(exclude=[...])
TensorPool.clear(exclude) 释放除 exclude 列表外的所有 tensor,同时重置引用计数:
def clear(self, exclude: list[TensorKey | TensorPoolKey] | None = None) -> None:
exclude_set = set(_normalize_key(k) for k in exclude) if exclude else set()
keys_to_remove = [k for k in self._store if k not in exclude_set]
for key in keys_to_remove:
self._store[key].release()
del self._store[key]
self._ref_counts.pop(key, None)---
ModelPool (`kdit/models/model_pool.py`)
- Owner: Executor
- 生命周期: 与 Executor 同生命周期,
clear_models()可手动清理 - 内容:
dict[ModelPoolKey, ModelBase](旧代码中存在ModelKey作为 key 的用法,已标记为 deprecated) - 用途: IONode 写入模型,InferNode 读取模型
---
DistributedGroupManager (`kdit/executor/distributed_group.py`)
- Owner: Executor
- 状态:
rank_id,world_size,_initialized - 用途: 提供
broadcast_tensors()能力,配合 tensor_pool 实现跨 rank 数据同步
---
相关文档
- Node 设计与 Pin 声明 → `node.md`
- PinHub 沙箱机制 → `pin-hub.md`
- Key 类型体系 → `../03_standards/key-system.md`
- 架构总览(Ownership 层级) → `overview.md`