From 6a4a70b02579974791b8afc4e50ea3e1dcdc5204 Mon Sep 17 00:00:00 2001 From: wangtao_ Date: Fri, 31 Jul 2026 10:28:26 +0800 Subject: [PATCH 1/2] feat(observability): fix model name - OlteRail: fix model name --- .../extensions/tracerotel/OtelRail.java | 16 ++++++++++++++++ 1 file changed, 16 insertions(+) diff --git a/src/main/java/com/openjiuwen/extensions/tracerotel/OtelRail.java b/src/main/java/com/openjiuwen/extensions/tracerotel/OtelRail.java index 286d5acbb..4e74a47fd 100644 --- a/src/main/java/com/openjiuwen/extensions/tracerotel/OtelRail.java +++ b/src/main/java/com/openjiuwen/extensions/tracerotel/OtelRail.java @@ -10,6 +10,7 @@ import com.openjiuwen.core.session.tracer.TraceAgentSpan; import com.openjiuwen.core.session.tracer.Tracer; import com.openjiuwen.core.session.tracer.TracerHandlerName; +import com.openjiuwen.core.foundation.llm.schema.ModelRequestConfig; import com.openjiuwen.core.singleagent.BaseAgent; import com.openjiuwen.core.singleagent.agents.ReActAgentConfig; import com.openjiuwen.core.singleagent.rail.AgentCallbackContext; @@ -422,6 +423,14 @@ private static String getModelName(AgentCallbackContext ctx) { /** * Extract model name from agent config object. * + *

Resolution order: + *

    + *
  1. {@code ReActAgentConfig.getModelName()} — explicitly set model name
  2. + *
  3. {@code ReActAgentConfig.getModelConfigObj().getModelName()} — model name from + * the request config (e.g. {@code "qwen-plus"})
  4. + *
  5. Empty {@link Optional} — caller falls back to {@code "LLM"}
  6. + *
+ * * @param config the agent config object (may be {@code null}) * @return an {@link Optional} containing the model name, or empty if not found * @since 0.1.7 @@ -435,6 +444,13 @@ private static Optional extractModelNameFromConfig(Object config) { if (name != null && !name.isEmpty()) { return Optional.of(name); } + ModelRequestConfig modelConfigObj = reactConfig.getModelConfigObj(); + if (modelConfigObj != null) { + String modelConfigName = modelConfigObj.getModelName(); + if (modelConfigName != null && !modelConfigName.isEmpty()) { + return Optional.of(modelConfigName); + } + } } return Optional.empty(); } From fde9fa2a25551edcee9dc7c346878f0e89361659 Mon Sep 17 00:00:00 2001 From: wangtao_ Date: Sat, 1 Aug 2026 11:53:46 +0800 Subject: [PATCH 2/2] =?UTF-8?q?fix(multiagent):=20=E4=BF=AE=E5=A4=8D?= =?UTF-8?q?=E5=88=86=E5=B1=82=E5=9B=A2=E9=98=9F=E5=B7=A5=E5=85=B7=E6=B3=A8?= =?UTF-8?q?=E5=86=8C=E5=86=B2=E7=AA=81=E5=8F=8A=E8=B5=84=E6=BA=90=E7=AE=A1?= =?UTF-8?q?=E7=90=86=E5=99=A8=E5=81=A5=E5=A3=AE=E6=80=A7=E9=97=AE=E9=A2=98?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ## 主要变更 ### 分层团队工具ID冲突修复 - HierarchicalDelegateTool/DelegateTool 的 toolId 添加 "_delegate_tool" 后缀, 避免与同名 agent 在 ResourceMgr tagMgr 中注册冲突 - 工具名称(toolName)保持不变,对 LLM 仍暴露子 agent ID ### 分层团队工具注入逻辑优化 - HierarchicalToolsTeam/HierarchicalMsgBusTeam 移除注入前的存在性检查 - 改用 addTool(tool, parentId, refresh=true) 强制刷新, 防止旧工具实例残留导致 "Tool instance not found" 错误 ### ResourceMgr 健壮性增强 - 批量注册(agent/workflow/tool)增加类型校验,非法类型抛出明确错误 - getTool 对空 ID 抛出 RESOURCE_ID_VALUE_INVALID 异常 - getResourceByTag 跳过 null 条目而非加入结果列表 - innerAddResource 重复注册时返回 Error 而非静默替换 - 新增 pythonTypeName 辅助方法,错误信息对齐 Python SDK ### TagMgr 空值处理对齐 Python - getResourcesTags 返回空列表替代 null - removeResourceTags 在 skipIfNotExists=false 时校验不存在的标签 - doRemoveResourceTags/doRemoveTag 清理空 tag 和空 resource 映射 ### ToolCard 模型类支持 - 新增 inputParamsRaw 字段及重载 builder 方法, 支持 Pydantic 模型类作为工具输入参数(对齐 Python input_params 类型) ### TaskManager 级联取消 - 实现子任务级联取消逻辑 ## 测试验证 - 4个分层团队冒烟测试全部通过 - 完整 smoke 套件 323 个测试 0 失败 0 错误 --- .../core/application/llm/LlmAgent.java | 53 +++++++++++ .../core/application/llm/LlmEventHandler.java | 4 +- .../core/common/task_manager/Task.java | 2 +- .../core/common/task_manager/TaskManager.java | 13 ++- .../core/common/utils/SchemaUtils.java | 3 +- .../InferenceAffinityModelClient.java | 6 +- .../OpenAiCompatibleModelClient.java | 6 +- .../core/foundation/tool/ToolCard.java | 70 ++++++++++++++ .../hierarchical_msgbus/DelegateTool.java | 6 +- .../HierarchicalMsgBusTeam.java | 13 +-- .../HierarchicalDelegateTool.java | 6 +- .../HierarchicalToolsTeam.java | 14 ++- .../core/runner/callback/CallbackChain.java | 12 +++ .../runner/callback/CallbackFramework.java | 4 +- .../runner/resourcemanager/ResourceMgr.java | 92 +++++++++++++++++-- .../core/runner/resourcemanager/TagMgr.java | 21 ++++- .../openjiuwen/core/workflow/Workflow.java | 3 +- 17 files changed, 292 insertions(+), 36 deletions(-) diff --git a/src/main/java/com/openjiuwen/core/application/llm/LlmAgent.java b/src/main/java/com/openjiuwen/core/application/llm/LlmAgent.java index a6480539b..c871a89b5 100644 --- a/src/main/java/com/openjiuwen/core/application/llm/LlmAgent.java +++ b/src/main/java/com/openjiuwen/core/application/llm/LlmAgent.java @@ -86,6 +86,59 @@ public LlmAgent(LlmAgentConfig agentConfig) { getController().setEventHandler(eventHandler); } + /** + * Add tools to this agent dynamically. + *

+ * Registers each tool in the ability manager and resource manager. + * Mirrors Python's LLMAgent.add_tools. + * + * @param tools the tools to add + * @since 0.1.7 + */ + public void addTools(List tools) { + if (tools == null || tools.isEmpty()) { + return; + } + String tag = agentConfig != null ? agentConfig.getId() : null; + for (Tool tool : tools) { + if (tool == null || tool.getCard() == null) { + continue; + } + registerPluginSchema(agentConfig, tool.getCard()); + this.getAbilityManager().add(tool.getCard()); + Runner.resourceMgr().addTool(tool, tag); + } + } + + /** + * Add workflows to this agent dynamically. + *

+ * Registers each workflow in the ability manager and resource manager. + * Mirrors Python's LLMAgent.add_workflows. + * + * @param workflows the workflows to add + * @since 0.1.7 + */ + public void addWorkflows(List workflows) { + if (workflows == null || workflows.isEmpty()) { + return; + } + String tag = agentConfig != null ? agentConfig.getId() : null; + for (Workflow workflow : workflows) { + if (workflow == null || workflow.getCard() == null) { + continue; + } + registerWorkflowSchema(agentConfig, workflow.getCard()); + this.getAbilityManager().add(workflow.getCard()); + WorkflowCard card = workflow.getCard(); + String workflowResourceId = WorkflowUtils.generateWorkflowKey(card.getId(), card.getVersion()); + WorkflowCard resourceCard = + WorkflowCard.builder().id(workflowResourceId).name(card.getName()).version(card.getVersion()) + .description(card.getDescription()).inputParams(card.getInputParams()).build(); + Runner.resourceMgr().addWorkflow(resourceCard, () -> workflow, tag); + } + } + /** * invoke. * diff --git a/src/main/java/com/openjiuwen/core/application/llm/LlmEventHandler.java b/src/main/java/com/openjiuwen/core/application/llm/LlmEventHandler.java index 4ce4e0ad2..fe8149fa0 100644 --- a/src/main/java/com/openjiuwen/core/application/llm/LlmEventHandler.java +++ b/src/main/java/com/openjiuwen/core/application/llm/LlmEventHandler.java @@ -810,8 +810,10 @@ private Model getModel() { } var modelInfo = agentConfig.getModel().modelInfo(); + String provider = agentConfig.getModel().modelProvider(); + Loggers.CONTROLLER.info("getModel: provider='{}', apiBase='{}'", provider, modelInfo.getApiBase()); ModelClientConfig clientConfig = ModelClientConfig.builder() - .clientProvider(agentConfig.getModel().modelProvider()).apiKey(modelInfo.getApiKey()) + .clientProvider(provider).apiKey(modelInfo.getApiKey()) .apiBase(modelInfo.getApiBase()).timeout(modelInfo.getTimeout()).verifySsl(modelInfo.isVerifySsl()) .sslCert(modelInfo.getSslCert()).headers(modelInfo.getHeaders()).build(); diff --git a/src/main/java/com/openjiuwen/core/common/task_manager/Task.java b/src/main/java/com/openjiuwen/core/common/task_manager/Task.java index 293cb99d5..52e2d1326 100644 --- a/src/main/java/com/openjiuwen/core/common/task_manager/Task.java +++ b/src/main/java/com/openjiuwen/core/common/task_manager/Task.java @@ -322,7 +322,7 @@ public boolean cancel(boolean isCascade, String reason, String cancelledBy) { markCancelled(); } if (isCascade && owner != null) { - owner.cascadeCancel(taskId, "parent_cancelled"); + owner.cascadeCancel(taskId, this.cancelReason); } if (timeoutHandle != null) { timeoutHandle.cancel(false); diff --git a/src/main/java/com/openjiuwen/core/common/task_manager/TaskManager.java b/src/main/java/com/openjiuwen/core/common/task_manager/TaskManager.java index cf52e11df..e67ceaa64 100644 --- a/src/main/java/com/openjiuwen/core/common/task_manager/TaskManager.java +++ b/src/main/java/com/openjiuwen/core/common/task_manager/TaskManager.java @@ -205,8 +205,17 @@ public Task createTask(Callable callable, String taskId, String name, Str */ public int cancelGroup(String group) { int count = 0; - for (Task task : registry.getByGroup(group)) { - if (task.cancel(false, "manual_cancel", null)) { + List groupTasks = registry.getByGroup(group); + for (Task task : groupTasks) { + // Skip tasks whose parent is also in the same group — they will be + // cancelled by the cascade originating from the root task. + boolean hasParentInGroup = task.getParentTaskId() != null + && registry.get(task.getParentTaskId()) != null + && group.equals(registry.get(task.getParentTaskId()).getGroup()); + if (hasParentInGroup) { + continue; + } + if (task.cancel(true, "manual_cancel", null)) { count++; } } diff --git a/src/main/java/com/openjiuwen/core/common/utils/SchemaUtils.java b/src/main/java/com/openjiuwen/core/common/utils/SchemaUtils.java index c218b6864..c1b323bd9 100644 --- a/src/main/java/com/openjiuwen/core/common/utils/SchemaUtils.java +++ b/src/main/java/com/openjiuwen/core/common/utils/SchemaUtils.java @@ -96,7 +96,8 @@ public static Map formatWithSchema(Map data, Map return dataWithDefaults; } catch (ValidationError e) { - throw e; + throw new ValidationError(StatusCode.SCHEMA_FORMAT_INVALID, null, null, e, + Map.of("reason", e.getMessage(), "data", String.valueOf(data))); } catch (Exception e) { throw new ValidationError(StatusCode.SCHEMA_FORMAT_INVALID, null, null, e, Map.of("reason", e.getMessage(), "data", String.valueOf(data))); diff --git a/src/main/java/com/openjiuwen/core/foundation/llm/model_clients/InferenceAffinityModelClient.java b/src/main/java/com/openjiuwen/core/foundation/llm/model_clients/InferenceAffinityModelClient.java index 5a188434c..cd22371ad 100644 --- a/src/main/java/com/openjiuwen/core/foundation/llm/model_clients/InferenceAffinityModelClient.java +++ b/src/main/java/com/openjiuwen/core/foundation/llm/model_clients/InferenceAffinityModelClient.java @@ -316,7 +316,11 @@ private HttpRequest buildJsonRequest(String suffix, Map body, Fl * @since 0.1.7 */ private String normalizedApiBase() { - return modelClientConfig.getApiBase().replaceAll("/+$", ""); + String base = modelClientConfig.getApiBase().replaceAll("/+$", ""); + if (base.endsWith("/v1/chat/completions")) { + base = base.substring(0, base.length() - "/v1/chat/completions".length()); + } + return base; } /** diff --git a/src/main/java/com/openjiuwen/core/foundation/llm/model_clients/OpenAiCompatibleModelClient.java b/src/main/java/com/openjiuwen/core/foundation/llm/model_clients/OpenAiCompatibleModelClient.java index 2ad156820..f5dd039df 100644 --- a/src/main/java/com/openjiuwen/core/foundation/llm/model_clients/OpenAiCompatibleModelClient.java +++ b/src/main/java/com/openjiuwen/core/foundation/llm/model_clients/OpenAiCompatibleModelClient.java @@ -360,7 +360,11 @@ private static void applyCallTimeout(Call call, Float timeoutOverride) { * @since 0.1.7 */ private String normalizedApiBase() { - return modelClientConfig.getApiBase().strip().replaceAll("/+$", ""); + String base = modelClientConfig.getApiBase().strip().replaceAll("/+$", ""); + if (base.endsWith("/chat/completions")) { + base = base.substring(0, base.length() - "/chat/completions".length()); + } + return base; } /** diff --git a/src/main/java/com/openjiuwen/core/foundation/tool/ToolCard.java b/src/main/java/com/openjiuwen/core/foundation/tool/ToolCard.java index 3b49d0b04..5b45fb119 100644 --- a/src/main/java/com/openjiuwen/core/foundation/tool/ToolCard.java +++ b/src/main/java/com/openjiuwen/core/foundation/tool/ToolCard.java @@ -30,6 +30,17 @@ public class ToolCard extends BaseCard { private Map inputParams = new HashMap<>(); + /** + * Raw inputParams object — can be a Map or a model Class. + *

+ * Mirrors Python's {@code input_params: Dict[str, Any] | Type[BaseModel]}. + * When set to a non-Map value (e.g., a Class), {@link #getInputParams()} returns + * the default empty Map, while {@link #getInputParamsRaw()} returns the original object. + * + * @since 0.1.14 + */ + private Object inputParamsRaw; + /** * Custom properties map. * @@ -68,6 +79,35 @@ public Map getInputParams() { */ public void setInputParams(Map inputParams) { this.inputParams = inputParams; + this.inputParamsRaw = inputParams; + } + + /** + * getInputParamsRaw. + *

+ * Returns the raw inputParams object, which may be a Map or a model Class. + * + * @return the raw inputParams object, or null if not set + * @since 0.1.14 + */ + public Object getInputParamsRaw() { + return inputParamsRaw; + } + + /** + * setInputParamsRaw. + *

+ * Sets the raw inputParams object. When the value is a Map, it also updates + * the typed {@code inputParams} field for backward compatibility. + * + * @param inputParamsRaw the raw inputParams object (Map or model Class) + * @since 0.1.14 + */ + public void setInputParamsRaw(Object inputParamsRaw) { + this.inputParamsRaw = inputParamsRaw; + if (inputParamsRaw instanceof Map) { + this.inputParams = (Map) inputParamsRaw; + } } /** @@ -113,6 +153,13 @@ public static class Builder extends BaseCard.Builder { */ protected Map inputParams = new HashMap<>(); + /** + * Raw inputParams for model class support. + * + * @since 0.1.14 + */ + protected Object inputParamsRaw; + /** * properties. * @@ -168,6 +215,26 @@ public Builder description(String description) { */ public Builder inputParams(Map inputParams) { this.inputParams = inputParams; + this.inputParamsRaw = inputParams; + return this; + } + + /** + * inputParams accepting any Object (Map or model Class). + *

+ * Mirrors Python's {@code input_params: Dict[str, Any] | Type[BaseModel]}. + * When a non-Map value is passed, it is stored in {@code inputParamsRaw} + * for retrieval via {@link ToolCard#getInputParamsRaw()}. + * + * @param inputParams the input parameters (Map or model Class) + * @return this builder + * @since 0.1.14 + */ + public Builder inputParams(Object inputParams) { + this.inputParamsRaw = inputParams; + if (inputParams instanceof Map) { + this.inputParams = (Map) inputParams; + } return this; } @@ -198,6 +265,9 @@ public ToolCard build() { card.setName(name); card.setDescription(description); card.setInputParams(inputParams); + if (inputParamsRaw != null) { + card.setInputParamsRaw(inputParamsRaw); + } card.setProperties(properties); return card; } diff --git a/src/main/java/com/openjiuwen/core/multiagent/teams/hierarchical_msgbus/DelegateTool.java b/src/main/java/com/openjiuwen/core/multiagent/teams/hierarchical_msgbus/DelegateTool.java index a397e1236..853dfe4a9 100644 --- a/src/main/java/com/openjiuwen/core/multiagent/teams/hierarchical_msgbus/DelegateTool.java +++ b/src/main/java/com/openjiuwen/core/multiagent/teams/hierarchical_msgbus/DelegateTool.java @@ -125,6 +125,7 @@ public Iterator stream(Map inputs, Map k */ private static ToolCard buildCard(String targetId, String targetDescription) { String toolName = targetId; + String toolId = targetId + "_delegate_tool"; String description = "Delegate a task to " + targetId + " for processing."; if (targetDescription != null && !targetDescription.isBlank()) { description += " " + targetDescription; @@ -138,6 +139,9 @@ private static ToolCard buildCard(String targetId, String targetDescription) { inputParams.put("type", "object"); inputParams.put("properties", properties); inputParams.put("required", List.of("message")); - return ToolCard.builder().id(toolName).name(toolName).description(description).inputParams(inputParams).build(); + // Use a distinct tool ID (suffixed with "_delegate_tool") to avoid + // ResourceMgr tagMgr collisions with the agent registered under the + // same ID. The tool name visible to the LLM remains the child agent ID. + return ToolCard.builder().id(toolId).name(toolName).description(description).inputParams(inputParams).build(); } } diff --git a/src/main/java/com/openjiuwen/core/multiagent/teams/hierarchical_msgbus/HierarchicalMsgBusTeam.java b/src/main/java/com/openjiuwen/core/multiagent/teams/hierarchical_msgbus/HierarchicalMsgBusTeam.java index 3e2ea26e2..f3aa6dc87 100644 --- a/src/main/java/com/openjiuwen/core/multiagent/teams/hierarchical_msgbus/HierarchicalMsgBusTeam.java +++ b/src/main/java/com/openjiuwen/core/multiagent/teams/hierarchical_msgbus/HierarchicalMsgBusTeam.java @@ -10,7 +10,6 @@ import com.openjiuwen.core.multiagent.BaseTeam; import com.openjiuwen.core.multiagent.schema.TeamCard; import com.openjiuwen.core.runner.Runner; -import com.openjiuwen.core.runner.base.TagMatchStrategy; import com.openjiuwen.core.session.AgentGroupSessionApi; import com.openjiuwen.core.singleagent.BaseAgent; import com.openjiuwen.core.singleagent.schema.AgentCard; @@ -113,18 +112,14 @@ private void injectDelegateTools(HierarchicalMsgBusTeamConfig config) { } for (String childId : children) { String toolName = childId; - if (supervisor.getAbilityManager().get(toolName) != null) { - continue; - } AgentCard childCard = getRuntime().getAgentCard(childId); String childDescription = childCard != null ? childCard.getDescription() : ""; DelegateTool tool = new DelegateTool(childId, childDescription, getRuntime(), supervisorId, teamId); supervisor.getAbilityManager().add(tool.getCard()); - Object existing = - Runner.resourceMgr().getTool(tool.getCard().getId(), supervisorId, TagMatchStrategy.ALL); - if (existing == null) { - Runner.resourceMgr().addTool(tool, supervisorId); - } + // Always register with refresh=true so stale tool instances from + // prior test runs are replaced. Without refresh, addTool silently + // skips when the toolId already exists, leaving the old instance. + Runner.resourceMgr().addTool(tool, supervisorId, true); Loggers.MULTI_AGENT.info("[HierarchicalMsgBusTeam:" + teamId + "] injected '" + toolName + "' -> '" + supervisorId + "'"); } diff --git a/src/main/java/com/openjiuwen/core/multiagent/teams/hierarchical_tools/HierarchicalDelegateTool.java b/src/main/java/com/openjiuwen/core/multiagent/teams/hierarchical_tools/HierarchicalDelegateTool.java index 2205fbc5d..f429f9ea6 100644 --- a/src/main/java/com/openjiuwen/core/multiagent/teams/hierarchical_tools/HierarchicalDelegateTool.java +++ b/src/main/java/com/openjiuwen/core/multiagent/teams/hierarchical_tools/HierarchicalDelegateTool.java @@ -137,6 +137,7 @@ public Iterator stream(Map inputs, Map k */ private static ToolCard buildCard(String targetId, AgentCard targetCard) { String toolName = targetId; + String toolId = targetId + "_delegate_tool"; String description = "Delegate a task to " + targetId + " for processing."; if (targetCard != null && targetCard.getDescription() != null && !targetCard.getDescription().isBlank()) { description = targetCard.getDescription(); @@ -151,7 +152,10 @@ private static ToolCard buildCard(String targetId, AgentCard targetCard) { } else { inputParams = defaultInputParams(); } - return ToolCard.builder().id(toolName).name(toolName).description(description).inputParams(inputParams).build(); + // Use a distinct tool ID (suffixed with "_delegate_tool") to avoid + // ResourceMgr tagMgr collisions with the agent registered under the + // same ID. The tool name visible to the LLM remains the child agent ID. + return ToolCard.builder().id(toolId).name(toolName).description(description).inputParams(inputParams).build(); } /** diff --git a/src/main/java/com/openjiuwen/core/multiagent/teams/hierarchical_tools/HierarchicalToolsTeam.java b/src/main/java/com/openjiuwen/core/multiagent/teams/hierarchical_tools/HierarchicalToolsTeam.java index 364ace138..beefdcc58 100644 --- a/src/main/java/com/openjiuwen/core/multiagent/teams/hierarchical_tools/HierarchicalToolsTeam.java +++ b/src/main/java/com/openjiuwen/core/multiagent/teams/hierarchical_tools/HierarchicalToolsTeam.java @@ -11,7 +11,6 @@ import com.openjiuwen.core.multiagent.schema.TeamCard; import com.openjiuwen.core.runner.Runner; import com.openjiuwen.core.runner.base.AgentProvider; -import com.openjiuwen.core.runner.base.TagMatchStrategy; import com.openjiuwen.core.session.AgentGroupSessionApi; import com.openjiuwen.core.session.stream.OutputSchema; import com.openjiuwen.core.singleagent.BaseAgent; @@ -208,9 +207,6 @@ private void injectChildCards(HierarchicalToolsTeamConfig config) { if (parent == null) { continue; } - if (parent.getAbilityManager().get(childId) != null) { - continue; - } AgentCard childCard = getRuntime().getAgentCard(childId); if (childCard == null) { Loggers.MULTI_AGENT.warning("[HierarchicalToolsTeam:" + teamId + "] skip tool injection for '" + childId @@ -220,10 +216,12 @@ private void injectChildCards(HierarchicalToolsTeamConfig config) { HierarchicalDelegateTool tool = new HierarchicalDelegateTool(childId, childCard, getRuntime(), parentId, teamId); parent.getAbilityManager().add(tool.getCard()); - Object existing = Runner.resourceMgr().getTool(tool.getCard().getId(), parentId, TagMatchStrategy.ALL); - if (existing == null) { - Runner.resourceMgr().addTool(tool, parentId); - } + // Always register with refresh=true so stale tool instances from + // prior test runs are replaced. Without refresh, addTool silently + // skips when the toolId already exists, leaving the old (possibly + // GC'd or wrong-scoped) instance — causing "Tool instance not found" + // errors at execution time. + Runner.resourceMgr().addTool(tool, parentId, true); Loggers.MULTI_AGENT.info("[HierarchicalToolsTeam:" + teamId + "] registered " + childId + " -> " + parentId + ".ability_manager"); } diff --git a/src/main/java/com/openjiuwen/core/runner/callback/CallbackChain.java b/src/main/java/com/openjiuwen/core/runner/callback/CallbackChain.java index ad34b5261..1c7a7aee8 100644 --- a/src/main/java/com/openjiuwen/core/runner/callback/CallbackChain.java +++ b/src/main/java/com/openjiuwen/core/runner/callback/CallbackChain.java @@ -95,6 +95,18 @@ public List getCallbacks() { return callbacks; } + /** + * Get the error handler registered for a specific callback, if any. + * + * @param callback the callback function + * @return the error handler, or {@code null} if none registered + * @since 0.1.7 + */ + public Function getErrorHandler( + Function, Object> callback) { + return errorHandlers.get(callback); + } + /** * Add callback to the chain. * diff --git a/src/main/java/com/openjiuwen/core/runner/callback/CallbackFramework.java b/src/main/java/com/openjiuwen/core/runner/callback/CallbackFramework.java index c720fbdca..7a1c05fca 100644 --- a/src/main/java/com/openjiuwen/core/runner/callback/CallbackFramework.java +++ b/src/main/java/com/openjiuwen/core/runner/callback/CallbackFramework.java @@ -617,7 +617,7 @@ public List triggerParallel(String event, Object[] args, Map triggerParallel(String event, Object[] args, Map> addAgents(List agents, Object tag) { validateTag(tag); } List> results = new ArrayList<>(); - for (AgentEntry entry : agents) { + List rawAgents = agents; + for (int i = 0; i < rawAgents.size(); i++) { + Object element = rawAgents.get(i); + if (!(element instanceof AgentEntry entry)) { + String gotType = pythonTypeName(element); + String length = element instanceof List l ? String.valueOf(l.size()) : "N/A"; + throw ErrorHelper.buildError(StatusCode.RESOURCE_PROVIDER_INVALID, "resource_type", "agent", + "reason", "invalid provider format at idx " + i + ": expected tuple[AgentCard, Callable]," + + " got " + gotType + " (length=" + length + ")"); + } results.add(innerAddResource(entry.card().getId(), entry.provider(), entry.card(), tag, "agent")); } return results; @@ -239,7 +249,16 @@ public List> addWorkflows(List workflows, Ob validateTag(tag); } List> results = new ArrayList<>(); - for (WorkflowEntry entry : workflows) { + List rawWorkflows = workflows; + for (int i = 0; i < rawWorkflows.size(); i++) { + Object element = rawWorkflows.get(i); + if (!(element instanceof WorkflowEntry entry)) { + String gotType = pythonTypeName(element); + String length = element instanceof List l ? String.valueOf(l.size()) : "N/A"; + throw ErrorHelper.buildError(StatusCode.RESOURCE_PROVIDER_INVALID, "resource_type", "workflow", + "reason", "invalid provider format at idx " + i + ": expected tuple[WorkflowCard, Callable]," + + " got " + gotType + " (length=" + length + ")"); + } results.add(innerAddResource(entry.card().getId(), entry.provider(), entry.card(), tag, "workflow")); } return results; @@ -353,8 +372,18 @@ public List> addTools(List tools, Object tag, boolean ref if (tag != null) { validateTag(tag); } + List rawTools = tools; + for (int i = 0; i < rawTools.size(); i++) { + Object element = rawTools.get(i); + if (!(element instanceof Tool)) { + String gotType = pythonTypeName(element); + throw ErrorHelper.buildError(StatusCode.RESOURCE_VALUE_INVALID, "resource_type", "tool", "reason", + "invalid tool type at index " + i + ": expected Tool, got " + gotType); + } + } List> results = new ArrayList<>(); - for (Tool tool : tools) { + for (int i = 0; i < rawTools.size(); i++) { + Tool tool = (Tool) rawTools.get(i); if (refresh) { refreshExistingResourceIfNeeded(tool.getCard().getId(), tag); } @@ -373,6 +402,10 @@ public List> addTools(List tools, Object tag, boolean ref * @since 0.1.7 */ public Object getTool(String toolId, Object tag, TagMatchStrategy tagMatchStrategy) { + if (toolId != null && (toolId.isEmpty() || toolId.isBlank())) { + throw ErrorHelper.buildError(StatusCode.RESOURCE_ID_VALUE_INVALID, "resource_type", "tool", "reason", + "id list cannot be empty or None"); + } return innerGetResources(toolId, tag, tagMatchStrategy, "tool"); } @@ -384,6 +417,10 @@ public Object getTool(String toolId, Object tag, TagMatchStrategy tagMatchStrate * @since 0.1.7 */ public Object getTool(String toolId) { + if (toolId == null || toolId.isEmpty() || toolId.isBlank()) { + throw ErrorHelper.buildError(StatusCode.RESOURCE_ID_VALUE_INVALID, "resource_type", "tool", "reason", + "id list cannot be empty or None"); + } return innerGetResources(toolId, null, TagMatchStrategy.ALL, "tool"); } @@ -1039,7 +1076,10 @@ public List getResourceByTag(String tag) { } List cards = new ArrayList<>(); for (String resourceId : resourceIds) { - cards.add(idToCard.get(resourceId)); + BaseCard card = idToCard.get(resourceId); + if (card != null) { + cards.add(card); + } } return cards; } @@ -1228,8 +1268,10 @@ private Result innerAddResource(String resourceId, Object resource, BaseC String resourceType) { try { if (tagMgr.hasResource(resourceId)) { - innerRemoveResources(resourceId, null, TagMatchStrategy.ALL, true, resourceType); - logger.info("replaced existing resource, id={}, type={}", resourceId, resourceType); + logger.info("resource already exist, id={}, type={}", resourceId, resourceType); + BaseError existError = ErrorHelper.buildError(StatusCode.RESOURCE_ADD_ERROR, + "card", resourceId, "reason", "resource already exist"); + return new Error<>(existError); } switch (resourceType) { case "workflow" -> resourceRegistry.workflow().addWorkflow(resourceId, (Supplier) resource); @@ -1271,6 +1313,10 @@ private List> innerRemoveResources(Object resourceId, Object tag, List idsToRemove; boolean isRemoveByTag = false; if (resourceId != null) { + if (resourceId instanceof String s && (s.isEmpty() || s.isBlank())) { + throw ErrorHelper.buildError(StatusCode.RESOURCE_ID_VALUE_INVALID, "resource_type", resourceType, + "reason", resourceType + " id list cannot be empty or None"); + } idsToRemove = normalizeIds(resourceId); } else { validateTag(tag); @@ -1624,6 +1670,40 @@ private static String getCardType(BaseCard card) { }; } + /** + * Map a Java object's type to the Python type name equivalent, so that error + * messages mirror the Python SDK output (e.g. {@code String -> "str"}). + * + * @param element the element whose Python type name is required + * @return the Python type name string + * @since 0.1.7 + */ + private static String pythonTypeName(Object element) { + if (element == null) { + return "NoneType"; + } + if (element instanceof String) { + return "str"; + } + if (element instanceof Integer || element instanceof Long || element instanceof Short + || element instanceof Byte) { + return "int"; + } + if (element instanceof Float || element instanceof Double) { + return "float"; + } + if (element instanceof Boolean) { + return "bool"; + } + if (element instanceof List) { + return "list"; + } + if (element instanceof Map) { + return "dict"; + } + return element.getClass().getSimpleName(); + } + @SuppressWarnings("unchecked") /** * normalizeIds. diff --git a/src/main/java/com/openjiuwen/core/runner/resourcemanager/TagMgr.java b/src/main/java/com/openjiuwen/core/runner/resourcemanager/TagMgr.java index c9f0eae29..89b25f659 100644 --- a/src/main/java/com/openjiuwen/core/runner/resourcemanager/TagMgr.java +++ b/src/main/java/com/openjiuwen/core/runner/resourcemanager/TagMgr.java @@ -156,7 +156,7 @@ public List getResourcesTags(String resourceId) { lock.lock(); try { Set tags = resourceTags.get(resourceId); - return tags != null ? new ArrayList<>(tags) : null; + return tags != null ? new ArrayList<>(tags) : Collections.emptyList(); } finally { lock.unlock(); } @@ -223,6 +223,16 @@ public List removeResourceTags(String resourceId, Object tags, boolean s throw ErrorHelper.buildError(StatusCode.RESOURCE_TAG_REMOVE_RESOURCE_TAG_ERROR, "resource_id", resourceId, "tags", String.valueOf(tags), "reason", "Resource does not exist"); } + Set currentTags = resourceTags.get(resourceId); + if (!skipIfNotExists) { + List nonExistent = tagsToRemove.stream() + .filter(t -> !currentTags.contains(t)) + .toList(); + if (!nonExistent.isEmpty()) { + throw ErrorHelper.buildError(StatusCode.RESOURCE_TAG_REMOVE_RESOURCE_TAG_ERROR, "resource_id", + resourceId, "tags", String.valueOf(nonExistent), "reason", "Tag does not exist"); + } + } return doRemoveResourceTags(resourceId, tagsToRemove); } finally { lock.unlock(); @@ -472,8 +482,14 @@ private List doRemoveResourceTags(String resourceId, List tagsTo Set res = tagToResource.get(tag); if (res != null) { res.remove(resourceId); + if (res.isEmpty() && !Tag.GLOBAL.equals(tag)) { + tagToResource.remove(tag); + } } } + if (currentTags.isEmpty()) { + resourceTags.remove(resourceId); + } return new ArrayList<>(currentTags); } @@ -492,6 +508,9 @@ private List doRemoveTag(String tag) { Set tags = resourceTags.get(resourceId); if (tags != null) { tags.remove(tag); + if (tags.isEmpty()) { + resourceTags.remove(resourceId); + } } affected.add(resourceId); } diff --git a/src/main/java/com/openjiuwen/core/workflow/Workflow.java b/src/main/java/com/openjiuwen/core/workflow/Workflow.java index 9cf919703..eace98fe4 100644 --- a/src/main/java/com/openjiuwen/core/workflow/Workflow.java +++ b/src/main/java/com/openjiuwen/core/workflow/Workflow.java @@ -910,7 +910,8 @@ private RuntimeException buildReceiveTimeoutError(boolean currentFirstFrame, lon formatTimeoutSeconds(receiveTimeoutMs), "reason", ""); } - if (configuredTimeoutMs > 0 && (executeTimeoutMs <= 0 || executeTimeoutMs > configuredTimeoutMs)) { + if (configuredTimeoutMs > 0 && (executeTimeoutMs <= 0 || executeTimeoutMs > configuredTimeoutMs) + && !deadlineReached) { return ErrorHelper.buildError(StatusCode.STREAM_OUTPUT_CHUNK_INTERVAL_TIMEOUT, "timeout", formatTimeoutSeconds(configuredTimeoutMs), "reason", ""); }