diff --git a/.agents/skills/llm-edgeflow-developer-guide/references/integration.md b/.agents/skills/llm-edgeflow-developer-guide/references/integration.md index 33576b84..664f1b34 100644 --- a/.agents/skills/llm-edgeflow-developer-guide/references/integration.md +++ b/.agents/skills/llm-edgeflow-developer-guide/references/integration.md @@ -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. diff --git a/doc/CHANGELOG.md b/doc/CHANGELOG.md index 674b9fed..db11d125 100644 --- a/doc/CHANGELOG.md +++ b/doc/CHANGELOG.md @@ -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/...`。 diff --git a/doc/dev_guide/business_onboarding.md b/doc/dev_guide/business_onboarding.md index c4581c2c..912ddcac 100644 --- a/doc/dev_guide/business_onboarding.md +++ b/doc/dev_guide/business_onboarding.md @@ -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、宿主指针和池内字符串均只在同步调用期间借用,不能保存。 diff --git a/include/adapter/io_binding.h b/include/adapter/io_binding.h index b36459dd..3afce178 100644 --- a/include/adapter/io_binding.h +++ b/include/adapter/io_binding.h @@ -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; }; diff --git a/include/adapter/io_binding_registry.h b/include/adapter/io_binding_registry.h index 65088df9..f8d6cb76 100644 --- a/include/adapter/io_binding_registry.h +++ b/include/adapter/io_binding_registry.h @@ -16,11 +16,6 @@ namespace llm_edgeflow { std::vector 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(); diff --git a/include/adapter/io_converter.h b/include/adapter/io_converter.h index 8dc8dd18..3105a936 100644 --- a/include/adapter/io_converter.h +++ b/include/adapter/io_converter.h @@ -204,8 +204,6 @@ struct InputConverterDefinition { std::string schema_id; std::vector external_slots; std::vector logical_ports; // 发布的内部逻辑输出端口 - // 可选的 Converter 专属上限;0 表示不设上限。 - size_t max_batch_size = 0; DecodeInputFn decode_fn = nullptr; }; @@ -220,8 +218,6 @@ struct OutputConverterDefinition { std::string schema_id; std::vector logical_ports; // 消费的内部逻辑输入端口 std::vector external_slots; - // 可选的 Converter 专属上限;0 表示不设上限。 - size_t max_batch_size = 0; EncodeOutputFn encode_fn = nullptr; }; diff --git a/include/nodes/control_authoring.h b/include/nodes/control_authoring.h index 6562028d..9f37979f 100644 --- a/include/nodes/control_authoring.h +++ b/include/nodes/control_authoring.h @@ -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 field_names, + FieldControlCommand(int cmd_id, std::string name, + std::vector 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& FieldNames() const noexcept { @@ -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; } @@ -146,31 +134,18 @@ class FieldControlCommand { [&](const ParamsT& current) -> NodeResult { ParamsT next = current; for (const auto& field_name : field_names_) { - if (strategy_ == ControlFieldStrategy::kReplace) { - if (!payload.contains(field_name)) { - return NodeResult::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::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::Failure( - NodeErrorKind::kBusinessError, assign_err, - node_error::control::kInvalidRequest); - } - } + if (!payload.contains(field_name)) { + return NodeResult::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::Failure( + NodeErrorKind::kBusinessError, assign_err, + node_error::control::kInvalidRequest); } } std::string val_err; @@ -188,7 +163,6 @@ class FieldControlCommand { } private: - ControlFieldStrategy strategy_; int cmd_id_; std::string name_; std::vector field_names_; @@ -199,16 +173,7 @@ class FieldControlCommand { inline FieldControlCommand ReplaceFields(int cmd_id, std::string name, std::vector 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 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)); } diff --git a/plans/FRAMEWORK_SIMPLIFICATION_PLAN.md b/plans/FRAMEWORK_SIMPLIFICATION_PLAN.md index d1389489..b8bb0a10 100644 --- a/plans/FRAMEWORK_SIMPLIFICATION_PLAN.md +++ b/plans/FRAMEWORK_SIMPLIFICATION_PLAN.md @@ -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. 目标与约束 @@ -42,7 +42,7 @@ | 进程重启后是否保留 | 不保留 | | 状态接口的形式 | 通用接口:既支持累积(如 embedding 历史),也支持切换(如某个模式一直保持到下一个特殊请求) | | BizDefinition 是否改为由 Binding 和转换器派生 | 不改,保留独立声明(原阶段 3b 取消,理由见第 9 节) | -| 阶段 4 的推进方式 | 先补齐设计和验收,评审确认后再实施 | +| 阶段 4 的推进方式 | 暂缓;有真实业务需求时,先按当时代码复核设计和验收,评审确认后再实施 | | 状态的提交边界 | 状态候选和输出发布对象全部准备成功后才最终提交;多个状态一起生效,或一起保持旧值 | | 状态的数据所有权 | 从产生起就是共享不可变对象;提交只转移或共享所有权;已有读取指针始终有效;所有读取路径对初始空状态的处理一致 | | 状态与请求数据如何配合 | 规定哪些端口可以直接读状态、哪些必须先广播到每条请求;新状态由写入节点生成,批内冲突指令由业务逻辑明确裁决 | @@ -83,7 +83,7 @@ 阶段 1、2 的设计与验收改由[业务开发者体验 RFC](DEVELOPER_EXPERIENCE_RFC.md)及其[详细设计](DEVELOPER_EXPERIENCE_DESIGN.md)维护。 -## 8. 阶段 4:跨请求状态(设计与验收,评审确认后实施) +## 8. 阶段 4:跨请求状态(暂缓;设计与验收留待真实需求出现时评审) ### 8.1 范围 diff --git a/src/adapter/biz/README.md b/src/adapter/biz/README.md index 9df38313..d39b10c4 100644 --- a/src/adapter/biz/README.md +++ b/src/adapter/biz/README.md @@ -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 读写集合推导全部业务契约。 diff --git a/src/adapter/deployment_preparation.cpp b/src/adapter/deployment_preparation.cpp index 9b2423d8..caed66a3 100644 --- a/src/adapter/deployment_preparation.cpp +++ b/src/adapter/deployment_preparation.cpp @@ -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; } diff --git a/src/adapter/io_binding_registry.cpp b/src/adapter/io_binding_registry.cpp index cedbf9e3..d35d271c 100644 --- a/src/adapter/io_binding_registry.cpp +++ b/src/adapter/io_binding_registry.cpp @@ -22,17 +22,6 @@ std::vector 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; @@ -206,11 +195,10 @@ bool IoBindingRegistry::Audit(std::vector* 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 绑定 diff --git a/src/adapter/io_catalog.cpp b/src/adapter/io_catalog.cpp index 388ef651..57591d68 100644 --- a/src/adapter/io_catalog.cpp +++ b/src/adapter/io_catalog.cpp @@ -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)}}; } @@ -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}}; diff --git a/src/cli/alg_pipeline_tool.cpp b/src/cli/alg_pipeline_tool.cpp index b54bf760..941aacf1 100644 --- a/src/cli/alg_pipeline_tool.cpp +++ b/src/cli/alg_pipeline_tool.cpp @@ -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}, diff --git a/src/core/node_definition_validation.cpp b/src/core/node_definition_validation.cpp index 482dbd3c..2439df34 100644 --- a/src/core/node_definition_validation.cpp +++ b/src/core/node_definition_validation.cpp @@ -21,8 +21,8 @@ const std::unordered_set& ValidProvenance() { } const std::unordered_set& ValidLifetimes() { - static const std::unordered_set kValidLifetimes = { - "request", "session", "global"}; + static const std::unordered_set kValidLifetimes = {"request", + "session"}; return kValidLifetimes; } diff --git a/src/core/pipeline_validator.cpp b/src/core/pipeline_validator.cpp index 2e13204e..83f4535c 100644 --- a/src/core/pipeline_validator.cpp +++ b/src/core/pipeline_validator.cpp @@ -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; } diff --git a/tests/integration/operator/test_operator_api.cpp b/tests/integration/operator/test_operator_api.cpp index e5c71e08..f52bdae3 100644 --- a/tests/integration/operator/test_operator_api.cpp +++ b/tests/integration/operator/test_operator_api.cpp @@ -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, diff --git a/tests/unit/adapter/test_adapter_purity.cpp b/tests/unit/adapter/test_adapter_purity.cpp index 73dc7290..e7cb1846 100644 --- a/tests/unit/adapter/test_adapter_purity.cpp +++ b/tests/unit/adapter/test_adapter_purity.cpp @@ -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, diff --git a/tests/unit/adapter/test_complex_converters.cpp b/tests/unit/adapter/test_complex_converters.cpp index 982613ba..973a1f07 100644 --- a/tests/unit/adapter/test_complex_converters.cpp +++ b/tests/unit/adapter/test_complex_converters.cpp @@ -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; } } diff --git a/tests/unit/adapter/test_io_binding_registry.cpp b/tests/unit/adapter/test_io_binding_registry.cpp index 8778831b..e2304d64 100644 --- a/tests/unit/adapter/test_io_binding_registry.cpp +++ b/tests/unit/adapter/test_io_binding_registry.cpp @@ -90,7 +90,6 @@ class IoBindingRegistryTest : public ::testing::Test { in_def.external_slots = { ExternalSlotDefinition("entity_in", "CompanyOperatorEntityInput", PortDirection::kInput, true, "entity_in")}; - in_def.max_batch_size = 64; in_def.logical_ports = { NodePortDefinition("input_sentences", "TextBatch", true, "1:1")}; in_def.decode_fn = &DummyDecode; @@ -103,7 +102,6 @@ class IoBindingRegistryTest : public ::testing::Test { out_def.external_slots = { ExternalSlotDefinition("entity_out", "CompanyOperatorEntityOutput", PortDirection::kOutput, true, "entity_out")}; - out_def.max_batch_size = 64; out_def.logical_ports = { NodePortDefinition("llm_answers", "TextBatch", true, "1:1")}; out_def.encode_fn = &DummyEncode; @@ -122,7 +120,8 @@ class IoBindingRegistryTest : public ::testing::Test { IoBindingRegistry::Instance().ClearForTesting(); } - void RegisterTestBizBinding(size_t max_batch_size = 0) { + void RegisterTestBizBinding( + size_t max_batch_size = kDefaultIoBindingMaxBatchSize) { IoBindingDefinition binding; binding.biz_name = "test_biz"; @@ -235,15 +234,6 @@ TEST_F(IoBindingRegistryTest, AuditRejectsUnregisteredConvertersAndBiz) { TEST_F(IoBindingRegistryTest, AuditAndPreparationRejectBindingWithoutBatchLimit) { - auto input = - *IoConverterRegistry::Instance().FindInputConverter("test.in.operator"); - auto output = - *IoConverterRegistry::Instance().FindOutputConverter("test.out.operator"); - IoConverterRegistry::Instance().ClearForTesting(); - input.max_batch_size = 0; - output.max_batch_size = 0; - ASSERT_TRUE(IoConverterRegistry::Instance().RegisterInputConverter(input)); - ASSERT_TRUE(IoConverterRegistry::Instance().RegisterOutputConverter(output)); RegisterTestBizBinding(0); std::vector errors; @@ -263,19 +253,10 @@ TEST_F(IoBindingRegistryTest, } TEST_F(IoBindingRegistryTest, BindingWithoutExplicitLimitUsesStandardDefault) { - auto input = - *IoConverterRegistry::Instance().FindInputConverter("test.in.operator"); - auto output = - *IoConverterRegistry::Instance().FindOutputConverter("test.out.operator"); - IoConverterRegistry::Instance().ClearForTesting(); - input.max_batch_size = 0; - output.max_batch_size = 0; - ASSERT_TRUE(IoConverterRegistry::Instance().RegisterInputConverter(input)); - ASSERT_TRUE(IoConverterRegistry::Instance().RegisterOutputConverter(output)); IoBindingDefinition binding; binding.biz_name = "test_biz"; - binding.input_converter_id = input.converter_id; - binding.output_converter_id = output.converter_id; + binding.input_converter_id = "test.in.operator"; + binding.output_converter_id = "test.out.operator"; ASSERT_TRUE(IoBindingRegistry::Instance().RegisterBinding(binding)); std::vector errors; EXPECT_TRUE(IoBindingRegistry::Instance().Audit(&errors)); @@ -1088,44 +1069,20 @@ TEST_F(IoBindingRegistryTest, SecondBindingForSameBizIsRejected) { EXPECT_EQ(prepared.output_specs.at("entity_out").type, "keyword_out"); } -TEST_F(IoBindingRegistryTest, EffectiveBatchLimitIncludesBindingBound) { - // 0 表示该来源不设上限;取最小的正值。 - struct Limits { - size_t binding; - size_t input; - size_t output; - size_t expected; - }; - const Limits cases[] = { - {1, 64, 64, 1}, {0, 64, 64, 64}, {0, 8, 16, 8}, - {128, 64, 64, 64}, {128, 8, 64, 8}, {128, 64, 16, 16}, - {64, 0, 0, 64}, {0, 0, 16, 16}, {32, 0, 64, 32}, - }; - auto input = - *IoConverterRegistry::Instance().FindInputConverter("test.in.operator"); - auto output = - *IoConverterRegistry::Instance().FindOutputConverter("test.out.operator"); +TEST_F(IoBindingRegistryTest, EffectiveBatchLimitIsBindingLimit) { const nlohmann::json document = { {"deployment", {{"io", {{"io_binding", "test_biz"}}}}}, {"pipeline", DefaultPipelineNodes()}}; - for (const auto& limits : cases) { - SCOPED_TRACE(::testing::Message() - << "binding=" << limits.binding << ", input=" << limits.input - << ", output=" << limits.output); + for (size_t limit : {1U, 32U, 128U}) { + SCOPED_TRACE(::testing::Message() << "binding=" << limit); IoBindingRegistry::Instance().ClearForTesting(); - IoConverterRegistry::Instance().ClearForTesting(); - input.max_batch_size = limits.input; - output.max_batch_size = limits.output; - ASSERT_TRUE(IoConverterRegistry::Instance().RegisterInputConverter(input)); - ASSERT_TRUE( - IoConverterRegistry::Instance().RegisterOutputConverter(output)); - RegisterTestBizBinding(limits.binding); + RegisterTestBizBinding(limit); PreparedDeployment prepared; DeploymentDiagnostic diagnostic; ASSERT_TRUE(PrepareDeploymentDocument(document, {}, &prepared, &diagnostic)) << diagnostic.message; - EXPECT_EQ(prepared.effective_max_batch_size, limits.expected); + EXPECT_EQ(prepared.effective_max_batch_size, limit); } } diff --git a/tests/unit/adapter/test_io_converters.cpp b/tests/unit/adapter/test_io_converters.cpp index bf8a4920..df60b3f9 100644 --- a/tests/unit/adapter/test_io_converters.cpp +++ b/tests/unit/adapter/test_io_converters.cpp @@ -80,7 +80,6 @@ TEST(IoConverterTest, RegisterAndFindInputConverter) { def.schema_id = "test_input"; def.external_slots = {ExternalSlotDefinition( "inputs", "int", PortDirection::kInput, true, "inputs")}; - def.max_batch_size = 64; def.logical_ports = {NodePortDefinition("texts", "TextBatch", true, "1:1")}; def.decode_fn = &DummyDecode; @@ -106,7 +105,6 @@ TEST(IoConverterTest, RegisterAndFindOutputConverter) { def.schema_id = "test_output"; def.external_slots = {ExternalSlotDefinition( "answers", "int", PortDirection::kOutput, true, "answers")}; - def.max_batch_size = 64; def.logical_ports = {NodePortDefinition("answers", "TextBatch", true, "1:1")}; def.encode_fn = &DummyEncode; @@ -137,7 +135,6 @@ TEST(IoConverterTest, RejectsInvalidDefinitions) { bad_in.schema_id = "test"; bad_in.external_slots = {ExternalSlotDefinition( "inputs", "int", PortDirection::kInput, true, "inputs")}; - bad_in.max_batch_size = 64; bad_in.logical_ports = { NodePortDefinition("texts", "TextBatch", true, "1:1")}; bad_in.decode_fn = nullptr; @@ -159,11 +156,10 @@ TEST(IoConverterTest, RejectsInvalidDefinitions) { bad_in.logical_ports.clear(); EXPECT_FALSE(IoConverterRegistry::Instance().RegisterInputConverter(bad_in)); - // max_batch_size == 0 表示转换器自身不设限,Binding 负责声明上限 - bad_in.converter_id = "unbounded.in"; + // 补齐必需字段后注册成功 + bad_in.converter_id = "complete.in"; bad_in.logical_ports = { NodePortDefinition("texts", "TextBatch", true, "1:1")}; - bad_in.max_batch_size = 0; EXPECT_TRUE(IoConverterRegistry::Instance().RegisterInputConverter(bad_in)); OutputConverterDefinition bad_out; @@ -172,7 +168,6 @@ TEST(IoConverterTest, RejectsInvalidDefinitions) { bad_out.schema_id = "test"; bad_out.external_slots = {ExternalSlotDefinition( "answers", "int", PortDirection::kOutput, true, "answers")}; - bad_out.max_batch_size = 64; bad_out.logical_ports = { NodePortDefinition("answers", "TextBatch", true, "1:1")}; bad_out.encode_fn = nullptr; @@ -243,7 +238,6 @@ TEST(IoConverterTest, bad_in.schema_id = "test_schema"; bad_in.external_slots = { ExternalSlotDefinition("slot1", "int", PortDirection::kInput, true, "")}; - bad_in.max_batch_size = 64; bad_in.logical_ports = { NodePortDefinition("texts", "TextBatch", true, "1:1")}; bad_in.decode_fn = &DummyDecode; @@ -257,7 +251,6 @@ TEST(IoConverterTest, bad_out.schema_id = "test_schema"; bad_out.external_slots = { ExternalSlotDefinition("slot1", "int", PortDirection::kOutput, true, "")}; - bad_out.max_batch_size = 64; bad_out.logical_ports = { NodePortDefinition("answers", "TextBatch", true, "1:1")}; bad_out.encode_fn = &DummyEncode; @@ -274,7 +267,6 @@ TEST(IoConverterTest, "int_suffix"), ExternalSlotDefinition("slot_second", "int", PortDirection::kInput, true, "int_suffix")}; - multi_in.max_batch_size = 64; multi_in.logical_ports = { NodePortDefinition("texts", "TextBatch", true, "1:1")}; multi_in.decode_fn = &DummyDecode; diff --git a/tests/unit/nodes/test_function_node.cpp b/tests/unit/nodes/test_function_node.cpp index fc3ba013..e2e0d654 100644 --- a/tests/unit/nodes/test_function_node.cpp +++ b/tests/unit/nodes/test_function_node.cpp @@ -204,7 +204,7 @@ auto MoveOnlyMapSpec() { static_assert(!std::is_copy_constructible_v); return MakeMapSpec(Input("input"), Output("output"), CleanConfig(), std::move(transform)) - .WithControls({PatchFields(3005, "patch_prefix", {"prefix"})}); + .WithControls({ReplaceFields(3005, "set_prefix", {"prefix"})}); } REGISTER_FUNCTION_NODE(MoveOnlyMapNode, MoveOnlyMapSpec()); @@ -664,7 +664,7 @@ struct ControlledMapParams { }; inline constexpr int kCmdReplaceMap = 3001; -inline constexpr int kCmdPatchMap = 3002; +inline constexpr int kCmdSuffixMap = 3002; inline std::string ControlledMapFn(const std::string& in, const ControlledMapParams& p) { @@ -698,8 +698,7 @@ inline auto ControlledMapSpec() { .WithControls({ ReplaceFields(kCmdReplaceMap, "replace_map", {"prefix", "multiplier"}), - PatchFields(kCmdPatchMap, "patch_map", - {"prefix", "suffix", "multiplier"}), + ReplaceFields(kCmdSuffixMap, "replace_suffix", {"suffix"}), }); } REGISTER_FUNCTION_NODE(ControlledMapNode, ControlledMapSpec()); @@ -2101,7 +2100,7 @@ TEST(ConfigurationSnapshotTest, MoveOnlyStateHandled) { // 函数式 Spec 的 WithControls 与 NodeHarness 测试 // --------------------------------------------------------------------------- -TEST(FunctionNodeTest, FunctionalMapSpecWithControlsReplaceAndPatch) { +TEST(FunctionNodeTest, FunctionalMapSpecWithControls) { NodeHarness harness("ControlledMapNode"); harness.Config( {{"prefix", "init_p:"}, {"suffix", ":init_s"}, {"multiplier", 1}}); @@ -2134,22 +2133,22 @@ TEST(FunctionNodeTest, FunctionalMapSpecWithControlsReplaceAndPatch) { EXPECT_EQ(res3.TextValues("output"), (std::vector{"rep_p:payloadpayload:init_s"})); - // 3. PatchFields 只带部分字段 -> 处理成功 - auto good_patch = harness.Control(kCmdPatchMap, R"({"suffix":":patch_s"})"); - EXPECT_EQ(good_patch.status, NodeControlStatus::kHandled); + // 3. 另一命令只控制 suffix -> 处理成功,其余字段不变 + auto good_suffix = harness.Control(kCmdSuffixMap, R"({"suffix":":patch_s"})"); + EXPECT_EQ(good_suffix.status, NodeControlStatus::kHandled); auto res4 = harness.Run(); ASSERT_TRUE(res4.ok()); EXPECT_EQ(res4.TextValues("output"), (std::vector{"rep_p:payloadpayload:patch_s"})); - // 4. PatchFields 传空对象 -> 拒绝 - auto empty_patch = harness.Control(kCmdPatchMap, R"({})"); - EXPECT_EQ(empty_patch.status, NodeControlStatus::kFailed); + // 4. 空对象缺少受控字段 -> 拒绝 + auto empty_payload = harness.Control(kCmdSuffixMap, R"({})"); + EXPECT_EQ(empty_payload.status, NodeControlStatus::kFailed); // 5. 语义校验失败 -> 拒绝并回滚状态 auto invalid_prefix = - harness.Control(kCmdPatchMap, R"({"prefix":"INVALID"})"); + harness.Control(kCmdReplaceMap, R"({"prefix":"INVALID","multiplier":2})"); EXPECT_EQ(invalid_prefix.status, NodeControlStatus::kFailed); EXPECT_NE(invalid_prefix.message.find("Invalid prefix disallowed"), std::string::npos); @@ -2494,11 +2493,11 @@ TEST(FunctionNodeTest, MixedControlDeclarationsRejectDuplicateIdInEitherOrder) { EXPECT_THROW( make_spec() .WithControl(custom, update) - .WithControls({PatchFields(4988, "patch_prefix", {"prefix"})}), + .WithControls({ReplaceFields(4988, "set_prefix", {"prefix"})}), std::invalid_argument); EXPECT_THROW( make_spec() - .WithControls({PatchFields(4988, "patch_prefix", {"prefix"})}) + .WithControls({ReplaceFields(4988, "set_prefix", {"prefix"})}) .WithControl(custom, update), std::invalid_argument); } @@ -2625,12 +2624,12 @@ TEST(FunctionNodeTest, RapidInterleavedControlsAndConcurrentProcesses) { } }); - // 写线程 2:快速 PatchFields + // 写线程 2:快速替换 suffix std::thread writer2([&]() { start.wait(); for (int i = 1; i <= 30; ++i) { std::string payload = "{\"suffix\":\":s" + std::to_string(i) + "\"}"; - auto res = node->Control(kCmdPatchMap, payload); + auto res = node->Control(kCmdSuffixMap, payload); EXPECT_EQ(res.status, NodeControlStatus::kHandled); std::this_thread::yield(); }