Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
Original file line number Diff line number Diff line change
Expand Up @@ -29,10 +29,9 @@ slot configurations, allocator and capacities. It is not a Pipeline Node and doe
6. Copy input data when the lifetime requires it, store request-scoped values in `AlgContext`, and pack output into leased pool slots only through the documented ownership contract.

Bindings default to the framework standard batch bound of 64; override
`IoBindingDefinition::max_batch_size` only when measurements require a smaller bound. Converters declare a limit only
when they have one of their own, and zero adds no bound. The effective limit is the smallest
positive value among the binding and its converters, and a binding where all three are zero fails
the registry audit and [deployment preparation](../../../../src/adapter/deployment_preparation.cpp).
`IoBindingDefinition::max_batch_size` only when measurements require a smaller bound. The binding is
the only source of this limit; converters do not declare one. A binding limit of zero fails the
registry audit and [deployment preparation](../../../../src/adapter/deployment_preparation.cpp).
[Operator creation](../../../../src/adapter/operator/operator_adapter.cpp) further caps the effective
Process batch limit at the output pool depth; a larger binding limit cannot relax another limit.

Expand Down
5 changes: 5 additions & 0 deletions doc/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,11 @@

## Unreleased

精简未使用的扩展点:批次上限只在 IoBinding 上声明,删除 Converter 的 `max_batch_size`(及
`EffectiveMaxBatchSize`),Catalog 的 Converter 不再导出该字段;字段 Control 只保留整体替换的
`ReplaceFields`,删除 `PatchFields` 及其策略枚举;端口存活期只接受 `request` 与 `session`,删除
未使用的 `global`;Catalog 与 `validate-io` 删除恒为 `operator` 的 `transport` 字段。

示例 Pipeline 配置、Demo 夹具与资产清单中的模型实例名去掉版本后缀(如 `embed_model_v1` /
`embed_model_v2` 改为 `embed_model`,`llm_model_v1` 改为 `llm_model`);Pipeline Studio 的 HTTP 接口
由 `/api/v1/...` 改为 `/api/...`。
Expand Down
6 changes: 2 additions & 4 deletions doc/dev_guide/business_onboarding.md
Original file line number Diff line number Diff line change
Expand Up @@ -130,10 +130,8 @@ JSON 请求是不同的输入约定。已有 Nodes 能完成算法,也不代
在 `IoBindingDefinition` 中指定 `biz_name`、`input_converter_id`、`output_converter_id`,
使用 `REGISTER_IO_BINDING` 注册;每个业务只注册一个绑定,Pipeline 以业务名选择它。
转换器的逻辑端口名就是 Blackboard Key,绑定不做改名;命名遵循本节后文的端口命名约定。
Binding 的批次上限默认为框架标准值 64,只有实测确需更小值时才覆盖 `max_batch_size`;
转换器只在自身确有限制时才声明上限,0 表示不设限。
有效上限取绑定与两个转换器中正值的最小值,Operator 再按实际输出池深收紧;
三者都为 0 时,注册审计和部署准备都会报错。
批次上限只在 Binding 上声明,默认为框架标准值 64,只有实测确需更小值时才覆盖 `max_batch_size`;
转换器不声明上限。Operator 再按实际输出池深收紧;上限为 0 时,注册审计和部署准备都会报错。

行函数返回的错误只需携带业务原因与字段路径,包装补充 converter 和样本位置。
输入行全部通过后才开始发布;输出 writer、宿主指针和池内字符串均只在同步调用期间借用,不能保存。
Expand Down
3 changes: 1 addition & 2 deletions include/adapter/io_binding.h
Original file line number Diff line number Diff line change
Expand Up @@ -18,8 +18,7 @@ struct IoBindingDefinition {

std::string input_converter_id;
std::string output_converter_id;
// 默认使用 Operator 标准上限;0 表示不设上限。有效上限取 binding 及其
// Converter 中最小的正值,且至少一方须为正。
// 单次 Process 的批大小上限,须为正;Operator 再按输出池深收紧。
size_t max_batch_size = kDefaultIoBindingMaxBatchSize;
};

Expand Down
5 changes: 0 additions & 5 deletions include/adapter/io_binding_registry.h
Original file line number Diff line number Diff line change
Expand Up @@ -16,11 +16,6 @@ namespace llm_edgeflow {
std::vector<std::string> EffectiveCapacityFields(
const ExternalSlotDefinition& slot);

// binding 及其 Converter 中最小的正上限;均未声明时为 0。
size_t EffectiveMaxBatchSize(const IoBindingDefinition& binding,
const InputConverterDefinition& input,
const OutputConverterDefinition& output);

class IoBindingRegistry {
public:
static IoBindingRegistry& Instance();
Expand Down
4 changes: 0 additions & 4 deletions include/adapter/io_converter.h
Original file line number Diff line number Diff line change
Expand Up @@ -204,8 +204,6 @@ struct InputConverterDefinition {
std::string schema_id;
std::vector<ExternalSlotDefinition> external_slots;
std::vector<NodePortDefinition> logical_ports; // 发布的内部逻辑输出端口
// 可选的 Converter 专属上限;0 表示不设上限。
size_t max_batch_size = 0;

DecodeInputFn decode_fn = nullptr;
};
Expand All @@ -220,8 +218,6 @@ struct OutputConverterDefinition {
std::string schema_id;
std::vector<NodePortDefinition> logical_ports; // 消费的内部逻辑输入端口
std::vector<ExternalSlotDefinition> external_slots;
// 可选的 Converter 专属上限;0 表示不设上限。
size_t max_batch_size = 0;

EncodeOutputFn encode_fn = nullptr;
};
Expand Down
73 changes: 19 additions & 54 deletions include/nodes/control_authoring.h
Original file line number Diff line number Diff line change
Expand Up @@ -20,23 +20,17 @@

namespace llm_edgeflow {

enum class ControlFieldStrategy {
kPatch,
kReplace,
};

// 字段 Control 命令:payload 必须给出全部受控字段,整体替换后再校验。
class FieldControlCommand {
public:
FieldControlCommand(ControlFieldStrategy strategy, int cmd_id,
std::string name, std::vector<std::string> field_names,
FieldControlCommand(int cmd_id, std::string name,
std::vector<std::string> field_names,
std::string description = "")
: strategy_(strategy),
cmd_id_(cmd_id),
: cmd_id_(cmd_id),
name_(std::move(name)),
field_names_(std::move(field_names)),
description_(std::move(description)) {}

ControlFieldStrategy Strategy() const noexcept { return strategy_; }
int Id() const noexcept { return cmd_id_; }
const std::string& Name() const noexcept { return name_; }
const std::vector<std::string>& FieldNames() const noexcept {
Expand Down Expand Up @@ -103,16 +97,10 @@ class FieldControlCommand {
prop["default"] = def.default_value;
}
props[fname] = std::move(prop);
if (strategy_ == ControlFieldStrategy::kReplace) {
req.push_back(fname);
}
req.push_back(fname);
}
schema["properties"] = std::move(props);
if (strategy_ == ControlFieldStrategy::kReplace) {
schema["required"] = std::move(req);
} else {
schema["minProperties"] = 1;
}
schema["required"] = std::move(req);
return schema;
}

Expand Down Expand Up @@ -146,31 +134,18 @@ class FieldControlCommand {
[&](const ParamsT& current) -> NodeResult<ParamsT> {
ParamsT next = current;
for (const auto& field_name : field_names_) {
if (strategy_ == ControlFieldStrategy::kReplace) {
if (!payload.contains(field_name)) {
return NodeResult<ParamsT>::Failure(
NodeErrorKind::kBusinessError,
"Missing required field in control payload: " +
field_name,
node_error::control::kInvalidRequest);
}
std::string assign_err;
if (!params.AssignField(field_name, payload[field_name], &next,
&assign_err)) {
return NodeResult<ParamsT>::Failure(
NodeErrorKind::kBusinessError, assign_err,
node_error::control::kInvalidRequest);
}
} else { // kPatch
if (payload.contains(field_name)) {
std::string assign_err;
if (!params.AssignField(field_name, payload[field_name],
&next, &assign_err)) {
return NodeResult<ParamsT>::Failure(
NodeErrorKind::kBusinessError, assign_err,
node_error::control::kInvalidRequest);
}
}
if (!payload.contains(field_name)) {
return NodeResult<ParamsT>::Failure(
NodeErrorKind::kBusinessError,
"Missing required field in control payload: " + field_name,
node_error::control::kInvalidRequest);
}
std::string assign_err;
if (!params.AssignField(field_name, payload[field_name], &next,
&assign_err)) {
return NodeResult<ParamsT>::Failure(
NodeErrorKind::kBusinessError, assign_err,
node_error::control::kInvalidRequest);
}
}
std::string val_err;
Expand All @@ -188,7 +163,6 @@ class FieldControlCommand {
}

private:
ControlFieldStrategy strategy_;
int cmd_id_;
std::string name_;
std::vector<std::string> field_names_;
Expand All @@ -199,16 +173,7 @@ class FieldControlCommand {
inline FieldControlCommand ReplaceFields(int cmd_id, std::string name,
std::vector<std::string> field_names,
std::string description = "") {
return FieldControlCommand(ControlFieldStrategy::kReplace, cmd_id,
std::move(name), std::move(field_names),
std::move(description));
}

inline FieldControlCommand PatchFields(int cmd_id, std::string name,
std::vector<std::string> field_names,
std::string description = "") {
return FieldControlCommand(ControlFieldStrategy::kPatch, cmd_id,
std::move(name), std::move(field_names),
return FieldControlCommand(cmd_id, std::move(name), std::move(field_names),
std::move(description));
}

Expand Down
6 changes: 3 additions & 3 deletions plans/FRAMEWORK_SIMPLIFICATION_PLAN.md
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@
| 1:业务源码自动收录 | 由 [业务开发者体验 RFC](DEVELOPER_EXPERIENCE_RFC.md) 的 WI-1 取代 |
| 2:Spec 签名提示 | 由 [业务开发者体验 RFC](DEVELOPER_EXPERIENCE_RFC.md) 的 WI-4 取代 |
| 3:Adapter 重复声明收敛 | 已完成,合入于 PR #150;现行规则见开发指南 |
| 4:跨请求状态 | 设计待评审,保留本文设计与验收 |
| 4:跨请求状态 | 暂缓:出现需要跨请求状态的真实业务前不实施;保留本文设计与验收,届时按当时代码复核后再评审 |
| 5:平台与调度参数归位 | 已完成,合入于 PR #150–#154 |

## 1. 目标与约束
Expand Down Expand Up @@ -42,7 +42,7 @@
| 进程重启后是否保留 | 不保留 |
| 状态接口的形式 | 通用接口:既支持累积(如 embedding 历史),也支持切换(如某个模式一直保持到下一个特殊请求) |
| BizDefinition 是否改为由 Binding 和转换器派生 | 不改,保留独立声明(原阶段 3b 取消,理由见第 9 节) |
| 阶段 4 的推进方式 | 先补齐设计和验收,评审确认后再实施 |
| 阶段 4 的推进方式 | 暂缓;有真实业务需求时,先按当时代码复核设计和验收,评审确认后再实施 |
| 状态的提交边界 | 状态候选和输出发布对象全部准备成功后才最终提交;多个状态一起生效,或一起保持旧值 |
| 状态的数据所有权 | 从产生起就是共享不可变对象;提交只转移或共享所有权;已有读取指针始终有效;所有读取路径对初始空状态的处理一致 |
| 状态与请求数据如何配合 | 规定哪些端口可以直接读状态、哪些必须先广播到每条请求;新状态由写入节点生成,批内冲突指令由业务逻辑明确裁决 |
Expand Down Expand Up @@ -83,7 +83,7 @@

阶段 1、2 的设计与验收改由[业务开发者体验 RFC](DEVELOPER_EXPERIENCE_RFC.md)及其[详细设计](DEVELOPER_EXPERIENCE_DESIGN.md)维护。

## 8. 阶段 4:跨请求状态(设计与验收,评审确认后实施)
## 8. 阶段 4:跨请求状态(暂缓;设计与验收留待真实需求出现时评审)

### 8.1 范围

Expand Down
2 changes: 1 addition & 1 deletion src/adapter/biz/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@

## 规范与契约
- 每个业务绑定文件通过 `BizDefinition` 声明业务名、Demo 名及完整 ingress/egress Blackboard 契约,并调用 `PipelineCatalog::RegisterBizDefinition` 登记。
- 显式声明该业务支持的 Operator 绑定,选择独立的输入/输出转换器,绑定批次上限默认为框架标准值 64,只有实测确需更小值时才覆盖 `max_batch_size`;转换器只在自身确有限制时才声明上限。
- 显式声明该业务支持的 Operator 绑定,选择独立的输入/输出转换器,绑定批次上限默认为框架标准值 64,只有实测确需更小值时才覆盖 `max_batch_size`;转换器不声明上限。
- 每个业务只注册一个绑定;业务配置在 Pipeline 的 `deployment.io.io_binding` 中填写业务名来选择它。
- 转换器的逻辑端口名即 Blackboard Key,绑定不做改名;复用转换器时沿用它声明的端口名,命名约定见[业务接入指南](../../../doc/dev_guide/business_onboarding.md#3-实现并注册转换器与绑定)。
- 保留独立转换器及完整业务 ingress/egress;不能仅根据 converter 读写集合推导全部业务契约。
7 changes: 3 additions & 4 deletions src/adapter/deployment_preparation.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -125,14 +125,13 @@ bool PrepareDeploymentDocument(const nlohmann::json& document,
return false;
}

const size_t max_batch = EffectiveMaxBatchSize(*binding, *in_conv, *out_conv);
const size_t max_batch = binding->max_batch_size;
if (max_batch == 0) {
if (diagnostic) {
diagnostic->code = "DEPLOYMENT_ERROR";
diagnostic->path = "/deployment/io/io_binding";
diagnostic->message =
"Binding '" + binding->biz_name +
"' declares no batch limit on the binding or its converters";
diagnostic->message = "Binding '" + binding->biz_name +
"' declares no batch limit (max_batch_size is 0)";
}
return false;
}
Expand Down
18 changes: 3 additions & 15 deletions src/adapter/io_binding_registry.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -22,17 +22,6 @@ std::vector<std::string> EffectiveCapacityFields(
return fields;
}

size_t EffectiveMaxBatchSize(const IoBindingDefinition& binding,
const InputConverterDefinition& input,
const OutputConverterDefinition& output) {
size_t limit = 0;
for (size_t candidate :
{binding.max_batch_size, input.max_batch_size, output.max_batch_size}) {
if (candidate != 0 && (limit == 0 || candidate < limit)) limit = candidate;
}
return limit;
}

IoBindingRegistry& IoBindingRegistry::Instance() {
static IoBindingRegistry instance;
return instance;
Expand Down Expand Up @@ -206,11 +195,10 @@ bool IoBindingRegistry::Audit(std::vector<std::string>* out_errors) const {
}
}

if (in_conv && out_conv &&
EffectiveMaxBatchSize(binding, *in_conv, *out_conv) == 0) {
if (binding.max_batch_size == 0) {
errors.push_back("Binding '" + biz_name +
"' declares no batch limit: set max_batch_size on the "
"binding or one of its converters");
"' declares no batch limit: max_batch_size must be "
"positive");
}

// 4. 检查对应槽位的 ValueType 绑定
Expand Down
5 changes: 0 additions & 5 deletions src/adapter/io_catalog.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -30,10 +30,8 @@ nlohmann::json InputConverterToJson(const InputConverterDefinition& conv) {
ports.push_back(PipelineCatalog::PortToJson(p.logical_name, p));

return {{"converter_id", conv.converter_id},
{"transport", "operator"},
{"schema_id", conv.schema_id},
{"external_type", ExternalType(conv.external_slots)},
{"max_batch_size", conv.max_batch_size},
{"external_slots", std::move(slots)},
{"logical_ports", std::move(ports)}};
}
Expand All @@ -46,17 +44,14 @@ nlohmann::json OutputConverterToJson(const OutputConverterDefinition& conv) {
ports.push_back(PipelineCatalog::PortToJson(p.logical_name, p));

return {{"converter_id", conv.converter_id},
{"transport", "operator"},
{"schema_id", conv.schema_id},
{"external_type", ExternalType(conv.external_slots)},
{"max_batch_size", conv.max_batch_size},
{"external_slots", std::move(slots)},
{"logical_ports", std::move(ports)}};
}

nlohmann::json IoBindingToJson(const IoBindingDefinition& b) {
return {{"biz_name", b.biz_name},
{"transport", "operator"},
{"input_converter_id", b.input_converter_id},
{"output_converter_id", b.output_converter_id},
{"max_batch_size", b.max_batch_size}};
Expand Down
1 change: 0 additions & 1 deletion src/cli/alg_pipeline_tool.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -637,7 +637,6 @@ int main(int argc, char* argv[]) {

nlohmann::json binding_info = {
{"biz_name", plan->binding.biz_name},
{"transport", "operator"},
{"input_converter_id", plan->binding.input_converter_id},
{"output_converter_id", plan->binding.output_converter_id},
{"effective_max_batch_size", plan->effective_max_batch_size},
Expand Down
4 changes: 2 additions & 2 deletions src/core/node_definition_validation.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -21,8 +21,8 @@ const std::unordered_set<std::string>& ValidProvenance() {
}

const std::unordered_set<std::string>& ValidLifetimes() {
static const std::unordered_set<std::string> kValidLifetimes = {
"request", "session", "global"};
static const std::unordered_set<std::string> kValidLifetimes = {"request",
"session"};
return kValidLifetimes;
}

Expand Down
1 change: 0 additions & 1 deletion src/core/pipeline_validator.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -166,7 +166,6 @@ bool ProvenanceCompatible(const std::string& producer,
int LifetimeRank(const std::string& lifetime) {
if (lifetime == "request") return 0;
if (lifetime == "session") return 1;
if (lifetime == "global") return 2;
return -1;
}

Expand Down
1 change: 0 additions & 1 deletion tests/integration/operator/test_operator_api.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -2426,7 +2426,6 @@ const bool g_reg_nested_output_components = []() {
OutputConverterDefinition odef;
odef.converter_id = "test_nested_output";

odef.max_batch_size = 64;
odef.schema_id = "test_nested_output";
odef.external_slots = {
ExternalSlotDefinition{"main", "test_nested_out", PortDirection::kOutput,
Expand Down
1 change: 0 additions & 1 deletion tests/unit/adapter/test_adapter_purity.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -1040,7 +1040,6 @@ TEST_F(AdapterPurityTest,
custom_in_def.external_slots = {
ExternalSlotDefinition("inputs", "CustomMultiFieldInput",
PortDirection::kInput, true, "custom_input")};
custom_in_def.max_batch_size = 64;
custom_in_def.logical_ports = {OutputPort(kInputSentences)};
custom_in_def.decode_fn = [](const ExternalInputBatchView& src,
const InputDecodeOptions& options,
Expand Down
5 changes: 1 addition & 4 deletions tests/unit/adapter/test_complex_converters.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -470,11 +470,8 @@ TEST_F(ComplexConvertersTest, AllEightBusinessesRegistered) {
converters.FindOutputConverter(binding->output_converter_id);
ASSERT_NE(input, nullptr) << biz;
ASSERT_NE(output, nullptr) << biz;
// 生产代码只在 binding 上声明一次上限。
// 生产 binding 使用框架标准批次上限。
EXPECT_EQ(binding->max_batch_size, 64U) << biz;
EXPECT_EQ(input->max_batch_size, 0U) << biz;
EXPECT_EQ(output->max_batch_size, 0U) << biz;
EXPECT_EQ(EffectiveMaxBatchSize(*binding, *input, *output), 64U) << biz;
}
}

Expand Down
Loading
Loading