From 0699bde4f8624832415974e0668d26bff8b19319 Mon Sep 17 00:00:00 2001 From: nicholascole Date: Thu, 9 Jul 2026 20:53:05 -0700 Subject: [PATCH 1/5] feat(server): switch embedded worker secrets to TaskDef.runtimeMetadata (target) Replace the interim __resolved_credentials__ enrich-script stamping with the target PR #1255 mechanism: when embedded, the worker's TaskDef declares its secret names on runtimeMetadata; the host resolves them at the SIMPLE task's own poll and injects values onto the wire-only Task.runtimeMetadata (never persisted, no JS injection, resolved at the right task). - AgentService.registerTaskDef(String, List) stamps runtimeMetadata only when EmbeddedMode.isEmbedded() && names non-empty; standalone leaves it empty (native execution-token pull still delivers). - Worker-tool branch passes AgentCompiler.collectToolCredentials(config) names. - Remove interim stamping: ToolCompiler workerCred helpers + enrich args, JavaScriptBuilder workerCredJson params/vars/injection, setWorkerCreds calls in AgentCompiler/MultiAgentCompiler. - Pin conductorVersion to the local runtimemeta build (has TaskDef.runtimeMetadata). - Replace ToolCompilerWorkerCredTest with WorkerRuntimeMetadataTest: asserts runtimeMetadata declared when embedded, empty standalone, and the enrich script no longer emits __resolved_credentials__ (validated fail-first). System-task delivery unchanged (LLM keys via host AI integration; HTTP/MCP/ planner headers via ${workflow.secrets}). Co-Authored-By: Claude Opus 4.8 (1M context) --- server/build.gradle | 10 +- .../runtime/compiler/AgentCompiler.java | 10 +- .../runtime/compiler/MultiAgentCompiler.java | 2 - .../runtime/compiler/ToolCompiler.java | 64 +------- .../runtime/service/AgentService.java | 18 ++- .../runtime/util/JavaScriptBuilder.java | 10 +- .../compiler/ToolCompilerWorkerCredTest.java | 126 ---------------- .../compiler/WorkerRuntimeMetadataTest.java | 139 ++++++++++++++++++ 8 files changed, 169 insertions(+), 210 deletions(-) delete mode 100644 server/conductor-agentspan/src/test/java/dev/agentspan/runtime/compiler/ToolCompilerWorkerCredTest.java create mode 100644 server/conductor-agentspan/src/test/java/dev/agentspan/runtime/compiler/WorkerRuntimeMetadataTest.java diff --git a/server/build.gradle b/server/build.gradle index ee3a3b9c..de1a31cc 100644 --- a/server/build.gradle +++ b/server/build.gradle @@ -15,11 +15,11 @@ repositories { // ── Version catalog ────────────────────────────────────────────── ext { - // AgentSpan compiles/tests against the published conductor. Embedded secret delivery uses - // ${workflow.secrets.NAME} references resolved at runtime by the HOST conductor's - // ParametersUtils.substituteSecrets / SecretsDAO (conductor-oss PR #1255) — agentspan itself - // references no PR #1255 API, so it does not need to build against it. - conductorVersion = '3.32.0-rc.3' + // TARGET branch: worker secrets use TaskDef.runtimeMetadata (conductor-oss PR #1255), so the + // server references TaskDef.setRuntimeMetadata and must build against a conductor that has it. + // Pinned to the local runtimemeta build (superset of 3.32.0-rc.3); revert to a published version + // once PR #1255 ships. (The interim on feature/embedded-secret-toggle builds against 3.32.0-rc.3.) + conductorVersion = '3.32.0-rc.3-runtimemeta-LOCAL' lombokVersion = '1.18.42' log4jVersion = '2.24.3' sqliteJdbcVersion = '3.47.0.0' diff --git a/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/compiler/AgentCompiler.java b/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/compiler/AgentCompiler.java index 52192239..513a364b 100644 --- a/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/compiler/AgentCompiler.java +++ b/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/compiler/AgentCompiler.java @@ -342,11 +342,11 @@ WorkflowDef compileSimple(AgentConfig config) { /** * Collect {@code toolName -> [credentialNames]} for the agent's tools: each tool's own declared - * credentials, falling back to the agent-level credential list. Fed to - * {@link ToolCompiler#setWorkerCreds} so SIMPLE worker tasks carry {@code __resolved_credentials__} - * secret references in embedded mode. + * credentials, falling back to the agent-level credential list. Used by {@code AgentService} to + * declare each worker tool's {@code TaskDef.runtimeMetadata} (embedded), so the host resolves the + * names at the SIMPLE task's poll and injects the values onto {@code Task.runtimeMetadata}. */ - static Map> collectToolCredentials(AgentConfig config) { + public static Map> collectToolCredentials(AgentConfig config) { List agentCreds = config.getCredentials() != null ? config.getCredentials() : List.of(); Map> map = new LinkedHashMap<>(); if (config.getTools() != null) { @@ -372,7 +372,6 @@ WorkflowDef compileWithTools(AgentConfig config) { List tools = config.getTools(); ToolCompiler tc = new ToolCompiler(); - tc.setWorkerCreds(collectToolCredentials(config)); boolean hasApproval = tools.stream().anyMatch(ToolConfig::isApprovalRequired); boolean hasMcp = tools.stream().anyMatch(t -> "mcp".equals(t.getToolType())); boolean hasApi = tools.stream().anyMatch(t -> "api".equals(t.getToolType())); @@ -747,7 +746,6 @@ WorkflowDef compileHybrid(AgentConfig config) { } ToolCompiler tc = new ToolCompiler(); - tc.setWorkerCreds(collectToolCredentials(config)); boolean hasApproval = allTools.stream().anyMatch(ToolConfig::isApprovalRequired); boolean hasMcp = allTools.stream().anyMatch(t -> "mcp".equals(t.getToolType())); boolean hasApi = allTools.stream().anyMatch(t -> "api".equals(t.getToolType())); diff --git a/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/compiler/MultiAgentCompiler.java b/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/compiler/MultiAgentCompiler.java index 724e505d..4ebf6ffa 100644 --- a/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/compiler/MultiAgentCompiler.java +++ b/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/compiler/MultiAgentCompiler.java @@ -1364,7 +1364,6 @@ WorkflowDef compileSwarmAgentWorkflow(AgentConfig agent, List transf allTools.addAll(transferTools); ToolCompiler tc = new ToolCompiler(); - tc.setWorkerCreds(AgentCompiler.collectToolCredentials(agent)); boolean hasApproval = allTools.stream().anyMatch(ToolConfig::isApprovalRequired); List> toolSpecs = tc.compileToolSpecs(allTools); @@ -1463,7 +1462,6 @@ private WorkflowDef compileSwarmAgentWorkflowWithSubAgents(AgentConfig agent, Li // 3. LLM step with transfer tools to decide whether to transfer to a peer ToolCompiler tc = new ToolCompiler(); - tc.setWorkerCreds(AgentCompiler.collectToolCredentials(agent)); List> transferToolSpecs = tc.compileToolSpecs(transferTools); WorkflowTask transferLlm = new WorkflowTask(); diff --git a/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/compiler/ToolCompiler.java b/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/compiler/ToolCompiler.java index 996d011d..80553608 100644 --- a/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/compiler/ToolCompiler.java +++ b/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/compiler/ToolCompiler.java @@ -119,55 +119,6 @@ private static Map escapeHeadersInConfig(Map cfg Map.entry("rag_search", "LLM_SEARCH_INDEX"), Map.entry("pull_workflow_messages", "PULL_WORKFLOW_MESSAGES")); - /** - * Per-tool credential names ({@code toolName -> [names]}) used to stamp - * {@code __resolved_credentials__} onto SIMPLE worker-tool tasks in EMBEDDED mode, so the host - * resolves each {@code ${workflow.secrets.NAME}} reference from its secret store at poll time. - * Set by {@link AgentCompiler} (preserves agent-level credential fallback); empty by default - * (no stamping — standalone, or non-worker tools whose secrets travel as headers). - */ - private Map> workerCreds = Map.of(); - - /** Inject per-tool credential names so worker-tool SIMPLE tasks can carry secret references. */ - void setWorkerCreds(Map> workerCreds) { - this.workerCreds = workerCreds != null ? workerCreds : Map.of(); - } - - /** True if a tool compiles to a SIMPLE worker task (executed by an external SDK worker). */ - private static boolean isWorkerTool(ToolConfig tool) { - String t = tool.getToolType() != null ? tool.getToolType() : "worker"; - return "SIMPLE".equals(TYPE_MAP.getOrDefault(t, "SIMPLE")); - } - - /** - * Build {@code {toolName -> {NAME: "${workflow.secrets.NAME}"}}} for this agent's SIMPLE - * worker tools, EMBEDDED only. The host resolves the references just-in-time at poll (via - * {@code ParametersUtils.substituteSecrets}); the SDK worker reads {@code __resolved_credentials__} - * from its task input and strips it. HTTP/MCP tools are excluded — their secrets travel as - * {@code ${workflow.secrets.NAME}} headers. - */ - private Map buildWorkerCredConfig(List tools) { - Map cfg = new LinkedHashMap<>(); - if (!EmbeddedMode.isEmbedded() || tools == null || workerCreds.isEmpty()) { - return cfg; - } - for (ToolConfig tool : tools) { - if (tool.getName() == null || !isWorkerTool(tool)) { - continue; - } - List names = workerCreds.get(tool.getName()); - if (names == null || names.isEmpty()) { - continue; - } - Map refs = new LinkedHashMap<>(); - for (String name : names) { - refs.put(name, "${workflow.secrets." + name + "}"); - } - cfg.put(tool.getName(), refs); - } - return cfg; - } - // ── Public API ─────────────────────────────────────────────────────── /** @@ -465,19 +416,9 @@ public Object[] buildEnrichTask(String agentName, String llmRef, List registered) for (ToolConfig tool : config.getTools()) { String tt = tool.getToolType(); if ("worker".equals(tt) && !registered.contains(tool.getName())) { - registerTaskDef(tool.getName()); + registerTaskDef( + tool.getName(), + AgentCompiler.collectToolCredentials(config).get(tool.getName())); registered.add(tool.getName()); } } @@ -1409,6 +1412,16 @@ private String extractSubagentIdentifier(Map event) { // ── Task registration ──────────────────────────────────────────── private void registerTaskDef(String taskName) { + registerTaskDef(taskName, null); + } + + /** + * Register a worker TaskDef. When embedded, {@code runtimeMetadata} declares the secret names the + * host must resolve at the SIMPLE task's poll and inject onto the wire-only + * {@code Task.runtimeMetadata} (conductor-oss PR #1255). Standalone leaves it empty — the native + * execution-token pull delivers secrets instead. + */ + private void registerTaskDef(String taskName, List runtimeMetadata) { TaskDef taskDef = new TaskDef(); taskDef.setName(taskName); taskDef.setRetryCount(2); @@ -1417,6 +1430,9 @@ private void registerTaskDef(String taskName) { taskDef.setTimeoutSeconds(0); taskDef.setResponseTimeoutSeconds(3600); taskDef.setTimeoutPolicy(TaskDef.TimeoutPolicy.RETRY); + if (EmbeddedMode.isEmbedded() && runtimeMetadata != null && !runtimeMetadata.isEmpty()) { + taskDef.setRuntimeMetadata(new ArrayList<>(runtimeMetadata)); + } try { TaskDef existing = metadataDAO.getTaskDef(taskName); diff --git a/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/util/JavaScriptBuilder.java b/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/util/JavaScriptBuilder.java index f7d07852..91557657 100644 --- a/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/util/JavaScriptBuilder.java +++ b/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/util/JavaScriptBuilder.java @@ -526,8 +526,7 @@ public static String enrichToolsScript( String cliConfigJson, String humanConfigJson, String wmqConfigJson, - String knownToolNamesJson, - String workerCredJson) { + String knownToolNamesJson) { return iife(" var httpCfg = " + httpConfigJson + ";" + " var mcpCfg = " + mcpConfigJson + ";" + " var mediaCfg = " + mediaConfigJson + ";" + " var agentToolCfg = " @@ -536,7 +535,6 @@ public static String enrichToolsScript( + cliConfigJson + ";" + " var humanCfg = " + humanConfigJson + ";" + " var wmqCfg = " + wmqConfigJson + ";" + " var knownNames = " + knownToolNamesJson + ";" - + " var workerCredCfg = " + workerCredJson + ";" + " var agentState = $.agentState || {};" + " var tcs = $.toolCalls || [];" + " var result = [];" @@ -697,7 +695,6 @@ public static String enrichToolsScript( + " if (t.type === 'SIMPLE') {" + " t.inputParameters._agent_state = agentState;" + " if ($.agentspanCtx) { t.inputParameters.__agentspan_ctx__ = $.agentspanCtx; }" - + " if (workerCredCfg[n]) { t.inputParameters.__resolved_credentials__ = workerCredCfg[n]; }" + " if (cliCfg[n]) { t.inputParameters._allowed_commands = cliCfg[n].allowedCommands; }" + " }" + " result.push(t);" @@ -1155,8 +1152,7 @@ public static String enrichToolsScriptDynamic( String ragConfigJson, String humanConfigJson, String wmqConfigJson, - String knownToolNamesJson, - String workerCredJson) { + String knownToolNamesJson) { return iife(" var httpCfg = " + httpConfigJson + ";" + " var mcpCfg = $.mcpConfig || {};" + " var apiCfg = $.apiConfig || {};" + " var mediaCfg = " @@ -1165,7 +1161,6 @@ public static String enrichToolsScriptDynamic( + ragConfigJson + ";" + " var humanCfg = " + humanConfigJson + ";" + " var wmqCfg = " + wmqConfigJson + ";" + " var knownNames = " + knownToolNamesJson + ";" - + " var workerCredCfg = " + workerCredJson + ";" + " var agentState = $.agentState || {};" + " var tcs = $.toolCalls || [];" + " var result = [];" @@ -1352,7 +1347,6 @@ public static String enrichToolsScriptDynamic( + " if (t.type === 'SIMPLE') {" + " t.inputParameters._agent_state = agentState;" + " if ($.agentspanCtx) { t.inputParameters.__agentspan_ctx__ = $.agentspanCtx; }" - + " if (workerCredCfg[n]) { t.inputParameters.__resolved_credentials__ = workerCredCfg[n]; }" + " }" + " result.push(t);" + " }" diff --git a/server/conductor-agentspan/src/test/java/dev/agentspan/runtime/compiler/ToolCompilerWorkerCredTest.java b/server/conductor-agentspan/src/test/java/dev/agentspan/runtime/compiler/ToolCompilerWorkerCredTest.java deleted file mode 100644 index 16be2b08..00000000 --- a/server/conductor-agentspan/src/test/java/dev/agentspan/runtime/compiler/ToolCompilerWorkerCredTest.java +++ /dev/null @@ -1,126 +0,0 @@ -/* - * Copyright (c) 2025 AgentSpan - * Licensed under the MIT License. - */ -package dev.agentspan.runtime.compiler; - -import static org.assertj.core.api.Assertions.assertThat; - -import java.util.List; -import java.util.Map; - -import org.graalvm.polyglot.Context; -import org.graalvm.polyglot.Value; -import org.junit.jupiter.api.AfterEach; -import org.junit.jupiter.api.Test; - -import com.fasterxml.jackson.databind.ObjectMapper; -import com.netflix.conductor.common.metadata.workflow.WorkflowTask; - -import dev.agentspan.runtime.model.AgentConfig; -import dev.agentspan.runtime.model.ToolConfig; -import dev.agentspan.runtime.util.EmbeddedMode; - -/** - * Verifies EMBEDDED mode stamps {@code __resolved_credentials__ = { NAME: "${workflow.secrets.NAME}" }} - * onto SIMPLE worker-tool tasks (via the enrich script), and that standalone / non-worker tools are - * left untouched. In embedded, the host resolves the references from its secret store at poll time. - */ -class ToolCompilerWorkerCredTest { - - private static final ObjectMapper MAPPER = new ObjectMapper(); - - @AfterEach - void resetEmbedded() { - new EmbeddedMode().setEmbedded(false); - } - - private static ToolConfig worker(String name, String... creds) { - return ToolConfig.builder() - .name(name) - .description(name) - .toolType("worker") - .config(Map.of("credentials", List.of(creds))) - .build(); - } - - private static ToolCompiler compilerFor(AgentConfig config) { - ToolCompiler tc = new ToolCompiler(); - tc.setWorkerCreds(AgentCompiler.collectToolCredentials(config)); - return tc; - } - - private static String enrichScript(ToolCompiler tc, List tools) { - Object[] r = tc.buildEnrichTask("agent", "agent_llm", tools, ""); - return (String) ((WorkflowTask) r[0]).getInputParameters().get("expression"); - } - - /** Execute the enrich script through GraalJS for one tool call; return that task's built map. */ - @SuppressWarnings("unchecked") - private static Map runEnrichForTool(String script, String toolName) throws Exception { - String wrapped = "var $ = {toolCalls: [{name: '" + toolName + "', taskReferenceName: 'call_1'," - + " inputParameters: {}}], agentState: {}, userPrompt: 'test'};" - + " JSON.stringify(" + script + ");"; - try (Context ctx = Context.newBuilder("js").allowAllAccess(true).build()) { - Value v = ctx.eval("js", wrapped); - Map outer = MAPPER.readValue(v.asString(), Map.class); - List> tasks = (List>) outer.get("dynamicTasks"); - return tasks.stream() - .filter(t -> toolName.equals(t.get("name"))) - .findFirst() - .orElseThrow(); - } - } - - @Test - void embedded_stampsPerToolSecretReference() { - new EmbeddedMode().setEmbedded(true); - ToolConfig gh = worker("gh", "GITHUB_TOKEN"); - AgentConfig config = AgentConfig.builder() - .name("a") - .model("openai/gpt-4o") - .tools(List.of(gh)) - .build(); - - String script = enrichScript(compilerFor(config), List.of(gh)); - - assertThat(script).contains("\"gh\":{\"GITHUB_TOKEN\":\"${workflow.secrets.GITHUB_TOKEN}\"}"); - } - - @Test - @SuppressWarnings("unchecked") - void embedded_injectsResolvedCredentialsOntoSimpleTask() throws Exception { - new EmbeddedMode().setEmbedded(true); - ToolConfig gh = worker("gh", "GITHUB_TOKEN"); - AgentConfig config = AgentConfig.builder() - .name("a") - .model("openai/gpt-4o") - .tools(List.of(gh)) - .build(); - - Map task = runEnrichForTool(enrichScript(compilerFor(config), List.of(gh)), "gh"); - - Map input = (Map) task.get("inputParameters"); - Map resolved = (Map) input.get("__resolved_credentials__"); - assertThat(resolved).containsEntry("GITHUB_TOKEN", "${workflow.secrets.GITHUB_TOKEN}"); - } - - @Test - @SuppressWarnings("unchecked") - void standalone_leavesWorkerTaskUntouched() throws Exception { - new EmbeddedMode().setEmbedded(false); - ToolConfig gh = worker("gh", "GITHUB_TOKEN"); - AgentConfig config = AgentConfig.builder() - .name("a") - .model("openai/gpt-4o") - .tools(List.of(gh)) - .build(); - - String script = enrichScript(compilerFor(config), List.of(gh)); - assertThat(script).doesNotContain("__resolved_credentials__\":{\"GITHUB_TOKEN"); - - Map task = runEnrichForTool(script, "gh"); - Map input = (Map) task.get("inputParameters"); - assertThat(input).doesNotContainKey("__resolved_credentials__"); - } -} diff --git a/server/conductor-agentspan/src/test/java/dev/agentspan/runtime/compiler/WorkerRuntimeMetadataTest.java b/server/conductor-agentspan/src/test/java/dev/agentspan/runtime/compiler/WorkerRuntimeMetadataTest.java new file mode 100644 index 00000000..267ff6a6 --- /dev/null +++ b/server/conductor-agentspan/src/test/java/dev/agentspan/runtime/compiler/WorkerRuntimeMetadataTest.java @@ -0,0 +1,139 @@ +/* + * Copyright (c) 2025 AgentSpan + * Licensed under the MIT License. + */ +package dev.agentspan.runtime.compiler; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import java.lang.reflect.Field; +import java.lang.reflect.Method; +import java.util.List; +import java.util.Map; + +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.Test; +import org.mockito.ArgumentCaptor; + +import com.netflix.conductor.common.metadata.tasks.TaskDef; +import com.netflix.conductor.common.metadata.workflow.WorkflowTask; +import com.netflix.conductor.dao.MetadataDAO; +import com.netflix.conductor.service.MetadataService; + +import dev.agentspan.runtime.model.AgentConfig; +import dev.agentspan.runtime.model.ToolConfig; +import dev.agentspan.runtime.service.AgentService; +import dev.agentspan.runtime.util.EmbeddedMode; + +/** + * Target-state worker-secret delivery (conductor-oss PR #1255): in EMBEDDED mode the worker's + * {@link TaskDef} declares its secret names on {@code runtimeMetadata}; the host resolves them at the + * SIMPLE task's own poll and injects the values onto the wire-only {@code Task.runtimeMetadata} — the + * enrich script never stamps {@code __resolved_credentials__} into persisted task input. Standalone + * leaves {@code runtimeMetadata} empty (the native execution-token pull delivers secrets instead). + */ +class WorkerRuntimeMetadataTest { + + @AfterEach + void resetEmbedded() { + new EmbeddedMode().setEmbedded(false); + } + + private static ToolConfig worker(String name, String... creds) { + return ToolConfig.builder() + .name(name) + .description(name) + .toolType("worker") + .config(Map.of("credentials", List.of(creds))) + .build(); + } + + private static AgentConfig agentWith(ToolConfig tool) { + return AgentConfig.builder() + .name("a") + .model("openai/gpt-4o") + .tools(List.of(tool)) + .build(); + } + + // ── The enrich script must NOT stamp __resolved_credentials__ (the retired interim) ── + + @Test + void enrichScript_neverStampsResolvedCredentials_embedded() { + new EmbeddedMode().setEmbedded(true); + ToolConfig gh = worker("gh", "GITHUB_TOKEN"); + Object[] r = new ToolCompiler().buildEnrichTask("agent", "agent_llm", List.of(gh), ""); + String script = (String) ((WorkflowTask) r[0]).getInputParameters().get("expression"); + assertThat(script).doesNotContain("__resolved_credentials__"); + } + + @Test + void enrichScript_neverStampsResolvedCredentials_standalone() { + new EmbeddedMode().setEmbedded(false); + ToolConfig gh = worker("gh", "GITHUB_TOKEN"); + Object[] r = new ToolCompiler().buildEnrichTask("agent", "agent_llm", List.of(gh), ""); + String script = (String) ((WorkflowTask) r[0]).getInputParameters().get("expression"); + assertThat(script).doesNotContain("__resolved_credentials__"); + } + + // ── AgentService declares runtimeMetadata on the worker TaskDef only when embedded ── + + @Test + void embedded_stampsRuntimeMetadataOnWorkerTaskDef() throws Exception { + new EmbeddedMode().setEmbedded(true); + TaskDef registered = registerWorkerTaskDef("gh", List.of("GITHUB_TOKEN")); + assertThat(registered.getRuntimeMetadata()).containsExactly("GITHUB_TOKEN"); + } + + @Test + void standalone_leavesRuntimeMetadataEmpty() throws Exception { + new EmbeddedMode().setEmbedded(false); + TaskDef registered = registerWorkerTaskDef("gh", List.of("GITHUB_TOKEN")); + assertThat(registered.getRuntimeMetadata()).isNullOrEmpty(); + } + + @Test + void collectToolCredentials_mapsWorkerToItsSecretNames() { + AgentConfig config = agentWith(worker("gh", "GITHUB_TOKEN", "GH_APP_ID")); + Map> creds = AgentCompiler.collectToolCredentials(config); + assertThat(creds.get("gh")).containsExactlyInAnyOrder("GITHUB_TOKEN", "GH_APP_ID"); + } + + /** + * Drive {@link AgentService}'s private {@code registerTaskDef(String, List)} with the credential + * names {@link AgentCompiler#collectToolCredentials} yields for {@code toolName}, and capture the + * {@link TaskDef} handed to {@code MetadataService.registerTaskDef}. + */ + private static TaskDef registerWorkerTaskDef(String toolName, List creds) throws Exception { + MetadataDAO metadataDAO = mock(MetadataDAO.class); + MetadataService metadataService = mock(MetadataService.class); + when(metadataDAO.getTaskDef(toolName)).thenReturn(null); + + AgentService service = new AgentService( + mock(dev.agentspan.runtime.compiler.AgentCompiler.class), + mock(dev.agentspan.runtime.normalizer.NormalizerRegistry.class), + mock(com.netflix.conductor.dao.ExecutionDAO.class), + metadataDAO, + mock(com.netflix.conductor.core.execution.WorkflowExecutor.class), + mock(com.netflix.conductor.service.WorkflowService.class), + mock(dev.agentspan.runtime.service.AgentStreamRegistry.class), + mock(com.netflix.conductor.service.ExecutionService.class), + mock(dev.agentspan.runtime.util.ProviderValidator.class)); + + Field msField = AgentService.class.getDeclaredField("metadataService"); + msField.setAccessible(true); + msField.set(service, metadataService); + + Method m = AgentService.class.getDeclaredMethod("registerTaskDef", String.class, List.class); + m.setAccessible(true); + m.invoke(service, toolName, creds); + + @SuppressWarnings("unchecked") + ArgumentCaptor> captor = ArgumentCaptor.forClass(List.class); + verify(metadataService).registerTaskDef(captor.capture()); + return captor.getValue().get(0); + } +} From 7ec9b4bec0b269fdd59d5fd6b8209e3acafd28b3 Mon Sep 17 00:00:00 2001 From: nicholascole Date: Thu, 9 Jul 2026 21:19:01 -0700 Subject: [PATCH 2/5] feat(sdk): read host-resolved worker secrets from Task.runtimeMetadata (target) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Switch all four SDK worker read-paths from the interim __resolved_credentials__ task-input key to the target Task.runtimeMetadata wire field (conductor-oss PR #1255): the host resolves the worker's declared TaskDef.runtimeMetadata secret names at the SIMPLE task's own poll and injects the values on the wire only — never persisted to task input. The native execution-token pull stays as the standalone fallback. - Java: WorkerManager.readRuntimeMetadata(task) reads task.getRuntimeMetadata(); dep bump conductor-client 5.0.1 -> 5.1.0 (+ mavenLocal for the local build); ReadResolvedCredentialsTest -> ReadRuntimeMetadataTest (fail-first validated). - Python: _dispatch reads task.runtime_metadata; test_resolved_credentials -> test_runtime_metadata (fail-first validated). - TypeScript: worker.ts reads task.runtimeMetadata (structural cast so it compiles against the current client until the OpenAPI type releases); drop the __resolved_credentials__ strip; new worker.test.ts host-delivered case (fail-first validated). credentials.test.ts accessor path already aligned. - C#: WorkerManager.ReadRuntimeMetadata(task) reads task.RuntimeMetadata; drop the input strip; RuntimeMetadataReadTests via reflection. Not built here (no dotnet). Client deps require releases carrying Task.runtimeMetadata; pins annotated to repin once those land. Ships together with the server switchover so a runtimeMetadata-declaring server is never paired with an SDK reading the old key. Co-Authored-By: Claude Opus 4.8 (1M context) --- .../src/Conductor.AI/Conductor.AI.csproj | 3 + sdk/csharp/src/Conductor.AI/WorkerManager.cs | 31 +++++----- .../RuntimeMetadataReadTests.cs | 55 ++++++++++++++++++ sdk/java/build.gradle | 7 ++- .../conductor/ai/internal/WorkerManager.java | 28 ++++----- .../conductor/ai/SerializerTest.java | 27 ++++++--- .../internal/ReadResolvedCredentialsTest.java | 53 ----------------- .../ai/internal/ReadRuntimeMetadataTest.java | 58 +++++++++++++++++++ sdk/python/pyproject.toml | 3 + .../conductor/ai/agents/runtime/_dispatch.py | 8 +-- ...redentials.py => test_runtime_metadata.py} | 14 +++-- sdk/typescript/src/worker.ts | 19 +++--- sdk/typescript/tests/unit/worker.test.ts | 33 +++++++++++ 13 files changed, 227 insertions(+), 112 deletions(-) create mode 100644 sdk/csharp/tests/Conductor.AI.Tests/RuntimeMetadataReadTests.cs delete mode 100644 sdk/java/src/test/java/org/conductoross/conductor/ai/internal/ReadResolvedCredentialsTest.java create mode 100644 sdk/java/src/test/java/org/conductoross/conductor/ai/internal/ReadRuntimeMetadataTest.java rename sdk/python/tests/unit/{test_resolved_credentials.py => test_runtime_metadata.py} (75%) diff --git a/sdk/csharp/src/Conductor.AI/Conductor.AI.csproj b/sdk/csharp/src/Conductor.AI/Conductor.AI.csproj index 3aff6011..87204666 100644 --- a/sdk/csharp/src/Conductor.AI/Conductor.AI.csproj +++ b/sdk/csharp/src/Conductor.AI/Conductor.AI.csproj @@ -38,6 +38,9 @@ + diff --git a/sdk/csharp/src/Conductor.AI/WorkerManager.cs b/sdk/csharp/src/Conductor.AI/WorkerManager.cs index 07c8c1dc..8d8a7aee 100644 --- a/sdk/csharp/src/Conductor.AI/WorkerManager.cs +++ b/sdk/csharp/src/Conductor.AI/WorkerManager.cs @@ -95,10 +95,9 @@ private async System.Threading.Tasks.Task ExecuteAsync(Task task, CancellationTo // Strip internal keys from the handler-visible input var handlerInput = inputData - .Where(kv => !string.Equals(kv.Key, "__agentspan_ctx__", StringComparison.OrdinalIgnoreCase) - && !string.Equals(kv.Key, "_agent_state", StringComparison.OrdinalIgnoreCase) - && !string.Equals(kv.Key, "__resolved_credentials__", StringComparison.OrdinalIgnoreCase) - && !string.Equals(kv.Key, "method", StringComparison.OrdinalIgnoreCase)) + .Where(kv => !string.Equals(kv.Key, "__agentspan_ctx__", StringComparison.OrdinalIgnoreCase) + && !string.Equals(kv.Key, "_agent_state", StringComparison.OrdinalIgnoreCase) + && !string.Equals(kv.Key, "method", StringComparison.OrdinalIgnoreCase)) .ToDictionary(kv => kv.Key, kv => kv.Value, StringComparer.OrdinalIgnoreCase); // Resolve and inject credentials via the centralized helper so the @@ -106,9 +105,10 @@ private async System.Threading.Tasks.Task ExecuteAsync(Task task, CancellationTo // process-wide lock. See docs/design/secret-injection-contract.md. // Tier-2 (env-injection) path; tier-1 (explicit-key) lands when the // user-facing API exposes a `credentials` parameter to agent factories. - // Embedded: the host resolves ${workflow.secrets.NAME} into __resolved_credentials__ - // at poll time. Prefer that map; otherwise fall back to the native token-pull. - var resolvedCredentials = ReadResolvedCredentials(inputData); + // Embedded: the host resolves the worker's declared TaskDef.runtimeMetadata secret names + // at poll time and delivers the values on the wire-only Task.RuntimeMetadata (never + // persisted). Prefer that map; otherwise fall back to the native token-pull. + var resolvedCredentials = ReadRuntimeMetadata(task); if (resolvedCredentials.Count == 0 && _credentialNames.Length > 0) { var creds = await _http.ResolveCredentialsAsync( @@ -221,18 +221,19 @@ or CredentialRateLimitException } /// - /// Read the host-delivered __resolved_credentials__ name→value map from task input - /// (embedded mode). The host resolves the stamped ${workflow.secrets.NAME} references at - /// poll time. Empty when absent (standalone → the native token-pull is used instead). + /// Read the host-delivered secret name→value map from Task.RuntimeMetadata (embedded + /// mode). The host resolves the worker's declared TaskDef.runtimeMetadata names from its + /// secret store at poll time and injects the values on the wire only — never persisted to task + /// input (conductor-oss PR #1255). Empty when absent (standalone → the native token-pull). /// - private static Dictionary ReadResolvedCredentials(Dictionary inputData) + private static Dictionary ReadRuntimeMetadata(Task task) { var result = new Dictionary(); - if (inputData.TryGetValue("__resolved_credentials__", out var rc) && rc.ValueKind == JsonValueKind.Object) + if (task?.RuntimeMetadata is { Count: > 0 } rm) { - foreach (var prop in rc.EnumerateObject()) - if (prop.Value.ValueKind == JsonValueKind.String) - result[prop.Name] = prop.Value.GetString()!; + foreach (var (k, v) in rm) + if (k is not null && v is not null) + result[k] = v; } return result; } diff --git a/sdk/csharp/tests/Conductor.AI.Tests/RuntimeMetadataReadTests.cs b/sdk/csharp/tests/Conductor.AI.Tests/RuntimeMetadataReadTests.cs new file mode 100644 index 00000000..7f999d42 --- /dev/null +++ b/sdk/csharp/tests/Conductor.AI.Tests/RuntimeMetadataReadTests.cs @@ -0,0 +1,55 @@ +// Copyright (c) 2025 Agentspan +// Licensed under the MIT License. + +using System.Collections.Generic; +using System.Reflection; +using Xunit; +using ModelTask = Conductor.Client.Models.Task; + +namespace Conductor.AI.Tests; + +/// +/// Embedded host-delivery read-path: the worker reads host-resolved secret values from +/// Task.RuntimeMetadata (wire-only, resolved by the host from the worker's declared +/// TaskDef.runtimeMetadata; conductor-oss PR #1255). Absent/empty yields an empty map +/// (standalone falls back to the native token-pull). +/// +public class RuntimeMetadataReadTests +{ + private static Dictionary Invoke(ModelTask task) + { + // WorkerPollLoop is internal; reach ReadRuntimeMetadata (private static) via reflection. + var type = typeof(CredentialScope).Assembly.GetType("Conductor.AI.WorkerPollLoop")!; + var method = type.GetMethod( + "ReadRuntimeMetadata", + BindingFlags.NonPublic | BindingFlags.Static)!; + return (Dictionary)method.Invoke(null, new object?[] { task })!; + } + + [Fact] + public void Extracts_host_delivered_values() + { + var task = new ModelTask( + taskId: "t1", + runtimeMetadata: new Dictionary + { + ["GITHUB_TOKEN"] = "ghp_host", + ["GH_APP_ID"] = "42", + }); + + var result = Invoke(task); + + Assert.Equal(2, result.Count); + Assert.Equal("ghp_host", result["GITHUB_TOKEN"]); + Assert.Equal("42", result["GH_APP_ID"]); + } + + [Fact] + public void Empty_when_absent_or_empty() + { + Assert.Empty(Invoke(new ModelTask(taskId: "t1"))); + Assert.Empty(Invoke(new ModelTask( + taskId: "t1", + runtimeMetadata: new Dictionary()))); + } +} diff --git a/sdk/java/build.gradle b/sdk/java/build.gradle index 23f927dd..af42140a 100644 --- a/sdk/java/build.gradle +++ b/sdk/java/build.gradle @@ -14,6 +14,9 @@ java { } repositories { + // TARGET: consumes the local conductor-client build that carries Task.runtimeMetadata + // (conductor-oss/java-sdk feat/task-runtime-metadata). Remove once that release lands. + mavenLocal() mavenCentral() } @@ -28,7 +31,9 @@ ext { // separately from the server engine (engine = 3.30.2); wire-compatible with // the 3.x task REST API, bundles the common DTOs, and provides native auth // via io.orkes.conductor.client.ApiClient (key/secret → token). - conductorClientVersion = '5.0.1' + // TARGET: 5.1.0 adds Task.runtimeMetadata (host-resolved worker secrets, wire-only). + // Currently a local mavenLocal build; repin to the published 5.1.0 once it releases. + conductorClientVersion = '5.1.0' } dependencies { diff --git a/sdk/java/src/main/java/org/conductoross/conductor/ai/internal/WorkerManager.java b/sdk/java/src/main/java/org/conductoross/conductor/ai/internal/WorkerManager.java index c1a6855e..2c6399ea 100644 --- a/sdk/java/src/main/java/org/conductoross/conductor/ai/internal/WorkerManager.java +++ b/sdk/java/src/main/java/org/conductoross/conductor/ai/internal/WorkerManager.java @@ -373,9 +373,10 @@ private TaskResult executeHandler(String taskName, Task task) { // problem. See docs/design/secret-injection-contract.md. Map resolvedSecrets = Collections.emptyMap(); List declared = taskCredentials.getOrDefault(taskName, Collections.emptyList()); - // Embedded: the host resolves ${workflow.secrets.NAME} into __resolved_credentials__ at - // poll time. Prefer that map; otherwise fall back to the native token-pull (standalone). - Map hostDelivered = readResolvedCredentials(inputData); + // Embedded: the host resolves the worker's declared TaskDef.runtimeMetadata secret names at + // poll time and delivers the values on the wire-only Task.runtimeMetadata (never persisted). + // Prefer that map; otherwise fall back to the native token-pull (standalone). + Map hostDelivered = readRuntimeMetadata(task); if (!hostDelivered.isEmpty()) { resolvedSecrets = hostDelivered; } else if (!declared.isEmpty()) { @@ -423,18 +424,19 @@ private TaskResult executeHandler(String taskName, Task task) { } /** - * Read the host-delivered {@code __resolved_credentials__} name→value map from task input - * (embedded mode). The host resolves the stamped {@code ${workflow.secrets.NAME}} references at - * poll time. Returns an empty map when absent (standalone → native token-pull is used instead). + * Read the host-delivered secret name→value map from {@code Task.runtimeMetadata} (embedded mode). + * The host resolves the worker's declared {@code TaskDef.runtimeMetadata} names from its secret + * store at poll time and injects the values on the wire only — never persisted to task input + * (conductor-oss PR #1255). Returns an empty map when absent (standalone → native token-pull). */ - private static Map readResolvedCredentials(Map inputData) { - if (inputData == null) return Collections.emptyMap(); - Object rc = inputData.get("__resolved_credentials__"); - if (!(rc instanceof Map m) || m.isEmpty()) return Collections.emptyMap(); + private static Map readRuntimeMetadata(Task task) { + if (task == null) return Collections.emptyMap(); + Map rm = task.getRuntimeMetadata(); + if (rm == null || rm.isEmpty()) return Collections.emptyMap(); Map out = new HashMap<>(); - for (Map.Entry e : m.entrySet()) { - if (e.getKey() != null && e.getValue() instanceof String s) { - out.put(e.getKey().toString(), s); + for (Map.Entry e : rm.entrySet()) { + if (e.getKey() != null && e.getValue() != null) { + out.put(e.getKey(), e.getValue()); } } return out; diff --git a/sdk/java/src/test/java/org/conductoross/conductor/ai/SerializerTest.java b/sdk/java/src/test/java/org/conductoross/conductor/ai/SerializerTest.java index 0623f5ab..6969c8d7 100644 --- a/sdk/java/src/test/java/org/conductoross/conductor/ai/SerializerTest.java +++ b/sdk/java/src/test/java/org/conductoross/conductor/ai/SerializerTest.java @@ -483,8 +483,10 @@ void llm_guardrail_requires_model_and_policy() { @Test @SuppressWarnings("unchecked") void on_condition_handoff_serialized_with_target() { - Agent supervisor = - Agent.builder().name("supervisor").model("anthropic/claude-sonnet-4-6").build(); + Agent supervisor = Agent.builder() + .name("supervisor") + .model("anthropic/claude-sonnet-4-6") + .build(); Agent worker = Agent.builder() .name("worker") .model("anthropic/claude-sonnet-4-6") @@ -847,8 +849,10 @@ void planner_context_emitted_with_text_and_url_entries() { // Mirrors the Python + TS serializer tests. The wire shape MUST be // byte-equal across SDKs so the server compiler sees the same // payload regardless of language. - Agent planner = - Agent.builder().name("planner_sub").model("anthropic/claude-sonnet-4-6").build(); + Agent planner = Agent.builder() + .name("planner_sub") + .model("anthropic/claude-sonnet-4-6") + .build(); ToolDef stub = ToolDef.builder() .name("stub") .description("stub") @@ -885,8 +889,10 @@ void planner_context_emitted_with_text_and_url_entries() { void planner_context_omitted_when_unset() { // Counterfactual: without plannerContext the field MUST NOT appear // on the wire. Pairs with the positive test — pins the gating. - Agent planner = - Agent.builder().name("planner_sub").model("anthropic/claude-sonnet-4-6").build(); + Agent planner = Agent.builder() + .name("planner_sub") + .model("anthropic/claude-sonnet-4-6") + .build(); ToolDef stub = ToolDef.builder() .name("stub") .description("stub") @@ -907,7 +913,8 @@ void planner_context_omitted_when_unset() { void planner_context_rejected_on_non_plan_execute_strategy() { // Same guard shape as planner=/fallback= — setting plannerContext // on anything other than PLAN_EXECUTE is a silent bug. - Agent sub = Agent.builder().name("sub").model("anthropic/claude-sonnet-4-6").build(); + Agent sub = + Agent.builder().name("sub").model("anthropic/claude-sonnet-4-6").build(); IllegalArgumentException e = assertThrows(IllegalArgumentException.class, () -> Agent.builder() .name("h") .model("anthropic/claude-sonnet-4-6") @@ -956,8 +963,10 @@ void parity_fields_serialized() { @Test void parity_fields_absent_when_unset() { - Agent agent = - Agent.builder().name("plain_agent").model("anthropic/claude-sonnet-4-6").build(); + Agent agent = Agent.builder() + .name("plain_agent") + .model("anthropic/claude-sonnet-4-6") + .build(); Map out = ser.serialize(agent); assertFalse(out.containsKey("reasoningEffort"), "reasoningEffort omitted when unset"); assertFalse(out.containsKey("maskedFields"), "maskedFields omitted when unset"); diff --git a/sdk/java/src/test/java/org/conductoross/conductor/ai/internal/ReadResolvedCredentialsTest.java b/sdk/java/src/test/java/org/conductoross/conductor/ai/internal/ReadResolvedCredentialsTest.java deleted file mode 100644 index 688bddee..00000000 --- a/sdk/java/src/test/java/org/conductoross/conductor/ai/internal/ReadResolvedCredentialsTest.java +++ /dev/null @@ -1,53 +0,0 @@ -/* - * Copyright (c) 2025 AgentSpan - * Licensed under the MIT License. - */ -package org.conductoross.conductor.ai.internal; - -import static org.junit.jupiter.api.Assertions.assertEquals; -import static org.junit.jupiter.api.Assertions.assertTrue; - -import java.lang.reflect.Method; -import java.util.HashMap; -import java.util.Map; - -import org.junit.jupiter.api.Test; - -/** - * Validates {@code WorkerManager.readResolvedCredentials} — the embedded host-delivery read-path - * that extracts {@code __resolved_credentials__} (resolved by the host from - * {@code ${workflow.secrets.NAME}}) from task input. Absent/empty → empty map (standalone falls - * back to the native token-pull). - */ -class ReadResolvedCredentialsTest { - - @SuppressWarnings("unchecked") - private static Map invoke(Map inputData) throws Exception { - Method m = WorkerManager.class.getDeclaredMethod("readResolvedCredentials", Map.class); - m.setAccessible(true); - return (Map) m.invoke(null, inputData); - } - - @Test - void extractsHostDeliveredStringValues() throws Exception { - Map rc = new HashMap<>(); - rc.put("GITHUB_TOKEN", "ghp_host"); - rc.put("NOT_A_STRING", 123); // non-string values are skipped - Map input = new HashMap<>(); - input.put("__resolved_credentials__", rc); - - Map out = invoke(input); - - assertEquals(1, out.size()); - assertEquals("ghp_host", out.get("GITHUB_TOKEN")); - } - - @Test - void emptyWhenKeyAbsentOrNull() throws Exception { - assertTrue(invoke(new HashMap<>()).isEmpty()); - assertTrue(invoke(null).isEmpty()); - Map emptyMap = new HashMap<>(); - emptyMap.put("__resolved_credentials__", new HashMap<>()); - assertTrue(invoke(emptyMap).isEmpty()); - } -} diff --git a/sdk/java/src/test/java/org/conductoross/conductor/ai/internal/ReadRuntimeMetadataTest.java b/sdk/java/src/test/java/org/conductoross/conductor/ai/internal/ReadRuntimeMetadataTest.java new file mode 100644 index 00000000..9c8d8fa1 --- /dev/null +++ b/sdk/java/src/test/java/org/conductoross/conductor/ai/internal/ReadRuntimeMetadataTest.java @@ -0,0 +1,58 @@ +/* + * Copyright (c) 2025 AgentSpan + * Licensed under the MIT License. + */ +package org.conductoross.conductor.ai.internal; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.lang.reflect.Method; +import java.util.HashMap; +import java.util.Map; + +import org.junit.jupiter.api.Test; + +import com.netflix.conductor.common.metadata.tasks.Task; + +/** + * Validates {@code WorkerManager.readRuntimeMetadata} — the embedded host-delivery read-path that + * extracts the host-resolved secret values from {@code Task.runtimeMetadata} (wire-only, resolved by + * the host from the worker's declared {@code TaskDef.runtimeMetadata}; conductor-oss PR #1255). + * Absent/empty → empty map (standalone falls back to the native token-pull). + */ +class ReadRuntimeMetadataTest { + + @SuppressWarnings("unchecked") + private static Map invoke(Task task) throws Exception { + Method m = WorkerManager.class.getDeclaredMethod("readRuntimeMetadata", Task.class); + m.setAccessible(true); + return (Map) m.invoke(null, task); + } + + private static Task taskWithRuntimeMetadata(Map rm) { + Task task = new Task(); + task.setRuntimeMetadata(rm); + return task; + } + + @Test + void extractsHostDeliveredValues() throws Exception { + Map rm = new HashMap<>(); + rm.put("GITHUB_TOKEN", "ghp_host"); + rm.put("GH_APP_ID", "42"); + + Map out = invoke(taskWithRuntimeMetadata(rm)); + + assertEquals(2, out.size()); + assertEquals("ghp_host", out.get("GITHUB_TOKEN")); + assertEquals("42", out.get("GH_APP_ID")); + } + + @Test + void emptyWhenAbsentOrEmpty() throws Exception { + assertTrue(invoke(null).isEmpty()); + assertTrue(invoke(new Task()).isEmpty()); + assertTrue(invoke(taskWithRuntimeMetadata(new HashMap<>())).isEmpty()); + } +} diff --git a/sdk/python/pyproject.toml b/sdk/python/pyproject.toml index c927fc84..66945664 100644 --- a/sdk/python/pyproject.toml +++ b/sdk/python/pyproject.toml @@ -13,6 +13,9 @@ license = {text = "MIT License"} # fall back to source builds and fail. Bump the ceiling once those wheels exist. requires-python = ">=3.10,<3.14" dependencies = [ + # TARGET: requires a conductor-python release that carries Task.runtime_metadata + # (host-resolved worker secrets, wire-only; conductor-oss/python-sdk feat/task-runtime-metadata). + # The field is additive; repin the floor to that release once it lands. "conductor-python>=1.3.11", "httpx>=0.24", "cloudpickle>=2.0", diff --git a/sdk/python/src/conductor/ai/agents/runtime/_dispatch.py b/sdk/python/src/conductor/ai/agents/runtime/_dispatch.py index 298c1f4d..a5188bb7 100644 --- a/sdk/python/src/conductor/ai/agents/runtime/_dispatch.py +++ b/sdk/python/src/conductor/ai/agents/runtime/_dispatch.py @@ -419,10 +419,10 @@ def tool_worker(task: Task) -> TaskResult: credential_names = list( _workflow_credentials.get(task.workflow_instance_id, []) ) - # Embedded: the host resolves ${workflow.secrets.NAME} into __resolved_credentials__ - # at poll time. Prefer that map; otherwise fall back to the native token-pull - # (standalone). Pop the key so it never leaks into the tool's kwargs. - host_delivered = task.input_data.pop("__resolved_credentials__", None) + # Embedded: the host resolves the worker's declared TaskDef.runtimeMetadata secret names + # at poll time and delivers the values on the wire-only Task.runtime_metadata (never + # persisted). Prefer that map; otherwise fall back to the native token-pull (standalone). + host_delivered = getattr(task, "runtime_metadata", None) resolved_secrets = {} if isinstance(host_delivered, dict) and host_delivered: resolved_secrets = { diff --git a/sdk/python/tests/unit/test_resolved_credentials.py b/sdk/python/tests/unit/test_runtime_metadata.py similarity index 75% rename from sdk/python/tests/unit/test_resolved_credentials.py rename to sdk/python/tests/unit/test_runtime_metadata.py index 48981770..6fe8e68c 100644 --- a/sdk/python/tests/unit/test_resolved_credentials.py +++ b/sdk/python/tests/unit/test_runtime_metadata.py @@ -1,6 +1,6 @@ -"""Embedded host-delivered credential path: the worker prefers -``__resolved_credentials__`` from task input (resolved by the host from -``${workflow.secrets.NAME}``) over the native execution-token pull. +"""Embedded host-delivered credential path: the worker prefers the host-resolved secret values on +``Task.runtime_metadata`` (wire-only, resolved by the host from the worker's declared +``TaskDef.runtimeMetadata``; conductor-oss PR #1255) over the native execution-token pull. """ from unittest.mock import patch @@ -20,10 +20,11 @@ def read_token() -> str: return make_tool_worker(td.func, td.name, tool_def=td) -def test_prefers_host_delivered_resolved_credentials(): +def test_prefers_host_delivered_runtime_metadata(): wrapper = _worker() task = Task() - task.input_data = {"__resolved_credentials__": {"GITHUB_TOKEN": "ghp_host_resolved"}} + task.input_data = {} + task.runtime_metadata = {"GITHUB_TOKEN": "ghp_host_resolved"} task.workflow_instance_id = "wf" task.task_id = "t" @@ -36,10 +37,11 @@ def test_prefers_host_delivered_resolved_credentials(): mock_fetcher.assert_not_called() -def test_falls_back_to_native_fetch_when_no_resolved_map(): +def test_falls_back_to_native_fetch_when_no_runtime_metadata(): wrapper = _worker() task = Task() task.input_data = {"__agentspan_ctx__": {"execution_token": "tok"}} + task.runtime_metadata = None task.workflow_instance_id = "wf" task.task_id = "t" diff --git a/sdk/typescript/src/worker.ts b/sdk/typescript/src/worker.ts index 8e92d159..d07e9a9a 100644 --- a/sdk/typescript/src/worker.ts +++ b/sdk/typescript/src/worker.ts @@ -212,7 +212,6 @@ export function stripInternalKeys(inputData: Record): Record - | undefined; + const hostDelivered = (task as { runtimeMetadata?: Record }) + .runtimeMetadata; let resolvedCredentials: Record = {}; if (hostDelivered && Object.keys(hostDelivered).length > 0) { diff --git a/sdk/typescript/tests/unit/worker.test.ts b/sdk/typescript/tests/unit/worker.test.ts index 55c99c98..65185da2 100644 --- a/sdk/typescript/tests/unit/worker.test.ts +++ b/sdk/typescript/tests/unit/worker.test.ts @@ -550,6 +550,39 @@ describe("WorkerManager", () => { expect(contextAvailable).toBe(true); }); + it("prefers host-delivered task.runtimeMetadata over the native pull (embedded)", async () => { + // Embedded: the host resolves the worker's declared TaskDef.runtimeMetadata secret names and + // delivers the values on the wire-only Task.runtimeMetadata. The worker must use that map and + // never hit the native /workers/secrets endpoint, even with no execution token present. + const serverUrl = "http://cred-embedded"; + const manager = new WorkerManager(serverUrl, {}, 100); + + let resolved: string | undefined; + manager.addWorker("rtm_task", async () => { + const { getCredential } = await import("../../src/credentials.js"); + resolved = await getCredential("MY_CRED"); + return { ok: true }; + }); + + const fetchSpy = vi.fn().mockResolvedValue({ ok: true, status: 200, text: async () => "" }); + vi.stubGlobal("fetch", fetchSpy); + + const wrapped = (manager as any)._wrapWorker((manager as any).pendingWorkers[0]); + await wrapped.execute({ + taskId: "task-1", + workflowInstanceId: "wf-1", + inputData: { arg1: "value" }, // no __agentspan_ctx__ execution token + runtimeMetadata: { MY_CRED: "host-value" }, + }); + + expect(resolved).toBe("host-value"); + expect( + fetchSpy.mock.calls.some( + ([u]: [unknown]) => typeof u === "string" && u.includes("/workers/secrets"), + ), + ).toBe(false); + }); + it("clears credential context after handler completes", async () => { const manager = new WorkerManager("http://test", {}, 100); From ba696062382dfc566a9c12e4dd5e0c892fe9631c Mon Sep 17 00:00:00 2001 From: nicholascole Date: Thu, 9 Jul 2026 21:51:27 -0700 Subject: [PATCH 3/5] fix(sdk): register worker TaskDefs create-only so embedded runtimeMetadata isn't clobbered MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The SDK self-registered each worker TaskDef with overwrite semantics, using a bare def. When embedded, the host server pre-registers the worker TaskDef and declares its secret names on TaskDef.runtimeMetadata (conductor-oss PR #1255) — overwriting with a bare def (the client TaskDef model carries no runtimeMetadata field) clobbered that and starved the host resolver, so Task.runtimeMetadata arrived empty and secrets never resolved. Fix, flag-free (no embedded prop): register create-only — create the TaskDef when absent, never overwrite one that exists. Embedded, the server's def (with runtimeMetadata) is left intact; standalone still gets the def created when missing. The existence check chooses correctly with no configuration, so it "just works" either way. - Python: ToolRegistry.register_tool_workers + the framework worker path use overwrite_task_def=False (conductor-python then does get_task_def → skip-if-exists → else register). - Java: WorkerManager.registerTaskDef checks metadataClient.getTaskDef first and skips when present. - Tests (fail-first validated): Python test_embedded_taskdef_registration asserts create-only; Java EmbeddedTaskDefRegistrationTest asserts no-overwrite-when-exists / create-when-absent. Surfaced by the local embedded webhook e2e. TS/C# SDKs don't self-register worker TaskDefs, so they were already correct. Co-Authored-By: Claude Opus 4.8 (1M context) --- .../conductor/ai/internal/WorkerManager.java | 18 ++++- .../EmbeddedTaskDefRegistrationTest.java | 69 +++++++++++++++++++ .../conductor/ai/agents/runtime/runtime.py | 4 +- .../ai/agents/runtime/tool_registry.py | 9 ++- .../test_embedded_taskdef_registration.py | 34 +++++++++ 5 files changed, 131 insertions(+), 3 deletions(-) create mode 100644 sdk/java/src/test/java/org/conductoross/conductor/ai/internal/EmbeddedTaskDefRegistrationTest.java create mode 100644 sdk/python/tests/unit/test_embedded_taskdef_registration.py diff --git a/sdk/java/src/main/java/org/conductoross/conductor/ai/internal/WorkerManager.java b/sdk/java/src/main/java/org/conductoross/conductor/ai/internal/WorkerManager.java index 2c6399ea..889d7f21 100644 --- a/sdk/java/src/main/java/org/conductoross/conductor/ai/internal/WorkerManager.java +++ b/sdk/java/src/main/java/org/conductoross/conductor/ai/internal/WorkerManager.java @@ -250,7 +250,23 @@ public void register( logger.info("Registered worker for task: {} (domain={})", taskName, domain); } - private void registerTaskDef(String taskName, int configuredTimeoutSeconds) { + /** + * Register the worker TaskDef create-only: create it when absent, but never overwrite one that + * already exists. When embedded, the host server pre-registers the worker TaskDef and declares + * its secret names on {@code TaskDef.runtimeMetadata} (conductor-oss PR #1255); overwriting here + * with a bare def (the client TaskDef model carries no runtimeMetadata) would clobber that and + * starve the host resolver. Standalone still gets the def created when absent. The existence + * check chooses correctly with no embedded flag. + */ + void registerTaskDef(String taskName, int configuredTimeoutSeconds) { + try { + if (metadataClient.getTaskDef(taskName) != null) { + logger.debug("Task def {} already exists — leaving it untouched (create-only)", taskName); + return; + } + } catch (Exception lookupFailed) { + // Not found (or lookup errored) — fall through and create it. + } try { long timeout = effectiveTaskTimeout(configuredTimeoutSeconds); TaskDef taskDef = new TaskDef(taskName); diff --git a/sdk/java/src/test/java/org/conductoross/conductor/ai/internal/EmbeddedTaskDefRegistrationTest.java b/sdk/java/src/test/java/org/conductoross/conductor/ai/internal/EmbeddedTaskDefRegistrationTest.java new file mode 100644 index 00000000..6979acf9 --- /dev/null +++ b/sdk/java/src/test/java/org/conductoross/conductor/ai/internal/EmbeddedTaskDefRegistrationTest.java @@ -0,0 +1,69 @@ +/* + * Copyright (c) 2025 AgentSpan + * Licensed under the MIT License. + */ +package org.conductoross.conductor.ai.internal; + +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.lang.reflect.Field; +import java.util.List; + +import org.conductoross.conductor.ai.AgentConfig; +import org.junit.jupiter.api.Test; + +import com.netflix.conductor.client.http.ConductorClient; +import com.netflix.conductor.client.http.MetadataClient; +import com.netflix.conductor.common.metadata.tasks.TaskDef; + +/** + * Worker TaskDefs are registered create-only: the SDK creates the def when absent but never + * overwrites one that already exists. When embedded, the host server pre-registers the worker + * TaskDef and declares its secret names on TaskDef.runtimeMetadata (conductor-oss PR #1255); + * overwriting here with a bare def (the client TaskDef model has no runtimeMetadata field) would + * clobber that and starve the host resolver. No embedded flag — the existence check decides. + */ +class EmbeddedTaskDefRegistrationTest { + + /** Fake client: reports whether a def "exists" and records any registration, without network. */ + private static final class RecordingMetadataClient extends MetadataClient { + private final boolean exists; + boolean registered = false; + + RecordingMetadataClient(boolean exists) { + this.exists = exists; + } + + @Override + public TaskDef getTaskDef(String taskType) { + return exists ? new TaskDef(taskType) : null; + } + + @Override + public void registerTaskDefs(List taskDefs) { + this.registered = true; + } + } + + private static boolean didRegister(boolean alreadyExists) throws Exception { + WorkerManager wm = new WorkerManager(new AgentConfig(), new ConductorClient()); + RecordingMetadataClient client = new RecordingMetadataClient(alreadyExists); + Field f = WorkerManager.class.getDeclaredField("metadataClient"); + f.setAccessible(true); + f.set(wm, client); + wm.registerTaskDef("check_secret", 300); + return client.registered; + } + + @Test + void doesNotOverwriteExistingTaskDef() throws Exception { + // Existing def (e.g. server-registered with runtimeMetadata) must be left untouched. + assertFalse(didRegister(true), "must not overwrite an existing TaskDef"); + } + + @Test + void createsTaskDefWhenAbsent() throws Exception { + assertTrue(didRegister(false), "must create the TaskDef when none exists"); + } +} diff --git a/sdk/python/src/conductor/ai/agents/runtime/runtime.py b/sdk/python/src/conductor/ai/agents/runtime/runtime.py index 486d2bbd..53529d9e 100644 --- a/sdk/python/src/conductor/ai/agents/runtime/runtime.py +++ b/sdk/python/src/conductor/ai/agents/runtime/runtime.py @@ -877,11 +877,13 @@ def prepare(self, agent: Any) -> None: _, workers = serialize_agent(agent) for w in workers: wrapper = make_tool_worker(w.func, w.name) + # Create-only: never overwrite an existing worker TaskDef (preserves server-set + # TaskDef.runtimeMetadata when embedded; see ToolRegistry.register_tool_workers). worker_task( task_definition_name=w.name, task_def=_default_task_def(w.name), register_task_def=True, - overwrite_task_def=True, + overwrite_task_def=False, lease_extend_enabled=True, )(wrapper) if workers: diff --git a/sdk/python/src/conductor/ai/agents/runtime/tool_registry.py b/sdk/python/src/conductor/ai/agents/runtime/tool_registry.py index 13f2d029..101a60f2 100644 --- a/sdk/python/src/conductor/ai/agents/runtime/tool_registry.py +++ b/sdk/python/src/conductor/ai/agents/runtime/tool_registry.py @@ -71,6 +71,13 @@ def register_tool_workers( if td.func is not None and td.tool_type in ("worker", "cli"): guardrails = td.guardrails if td.guardrails else None wrapper = make_tool_worker(td.func, td.name, guardrails=guardrails, tool_def=td) + # Create-only (overwrite_task_def=False): register the TaskDef if it does not exist, + # but never overwrite one that does. When embedded, the host server pre-registers the + # worker TaskDef and declares its secret names on TaskDef.runtimeMetadata (conductor-oss + # PR #1255); overwriting here with a bare def (the client TaskDef model carries no + # runtimeMetadata) would clobber that and starve the host resolver. Standalone still + # gets the def created when absent. No embedded flag needed — the existence check makes + # the right choice automatically. worker_task( task_definition_name=td.name, task_def=_default_task_def( @@ -80,7 +87,7 @@ def register_tool_workers( retry_policy=td.retry_policy, ), register_task_def=True, - overwrite_task_def=True, + overwrite_task_def=False, domain=domain if (agent_stateful or td.stateful) else None, lease_extend_enabled=True, )(wrapper) diff --git a/sdk/python/tests/unit/test_embedded_taskdef_registration.py b/sdk/python/tests/unit/test_embedded_taskdef_registration.py new file mode 100644 index 00000000..d28d4dcf --- /dev/null +++ b/sdk/python/tests/unit/test_embedded_taskdef_registration.py @@ -0,0 +1,34 @@ +"""Worker TaskDefs are registered create-only (overwrite_task_def=False): the SDK creates the def +when absent but never overwrites an existing one. When embedded, the host server pre-registers the +worker TaskDef and declares its secret names on TaskDef.runtimeMetadata (conductor-oss PR #1255); +overwriting here with a bare def (the client TaskDef model has no runtimeMetadata field) would clobber +that and starve the host resolver. This needs no embedded flag — the existence check chooses correctly. +""" + +from unittest.mock import patch + +from conductor.ai.agents.runtime.tool_registry import ToolRegistry +from conductor.ai.agents.tool import tool + + +def _worker_task_kwargs(): + @tool(credentials=["DEMO_SECRET"]) + def check_secret() -> dict: + return {"ok": True} + + calls = [] + + def fake_worker_task(**kwargs): + calls.append(kwargs) + return lambda fn: fn # decorator passthrough + + with patch("conductor.client.worker.worker_task.worker_task", side_effect=fake_worker_task): + ToolRegistry().register_tool_workers([check_secret], "secret_agent") + return next(c for c in calls if c.get("task_definition_name") == "check_secret") + + +def test_worker_taskdef_is_create_only_never_overwrite(): + kwargs = _worker_task_kwargs() + # create-only: register when missing, but never overwrite (preserves server runtimeMetadata). + assert kwargs["register_task_def"] is True + assert kwargs["overwrite_task_def"] is False From eaff0469aee048eead73d93be9d970bc472c3099 Mon Sep 17 00:00:00 2001 From: Viren Baraiya Date: Fri, 10 Jul 2026 00:44:40 -0700 Subject: [PATCH 4/5] updates --- server/build.gradle | 2 +- .../runtime/util/EnrichToolsScriptTest.java | 2 +- .../runtime/compiler/AgentCompiler.java | 20 +++- .../runtime/service/AgentService.java | 31 ++--- .../compiler/WorkerRuntimeMetadataTest.java | 107 ++++++++++++++++++ 5 files changed, 146 insertions(+), 16 deletions(-) diff --git a/server/build.gradle b/server/build.gradle index de1a31cc..abfabb53 100644 --- a/server/build.gradle +++ b/server/build.gradle @@ -19,7 +19,7 @@ ext { // server references TaskDef.setRuntimeMetadata and must build against a conductor that has it. // Pinned to the local runtimemeta build (superset of 3.32.0-rc.3); revert to a published version // once PR #1255 ships. (The interim on feature/embedded-secret-toggle builds against 3.32.0-rc.3.) - conductorVersion = '3.32.0-rc.3-runtimemeta-LOCAL' + conductorVersion = '3.32.0-rc.5' lombokVersion = '1.18.42' log4jVersion = '2.24.3' sqliteJdbcVersion = '3.47.0.0' diff --git a/server/conductor-agentspan-server/src/test/java/dev/agentspan/runtime/util/EnrichToolsScriptTest.java b/server/conductor-agentspan-server/src/test/java/dev/agentspan/runtime/util/EnrichToolsScriptTest.java index 6949246d..3f97bfbc 100644 --- a/server/conductor-agentspan-server/src/test/java/dev/agentspan/runtime/util/EnrichToolsScriptTest.java +++ b/server/conductor-agentspan-server/src/test/java/dev/agentspan/runtime/util/EnrichToolsScriptTest.java @@ -55,7 +55,7 @@ private List> enrichWithAgentTools( private List> enrichWithConfigs( String httpJson, String agentToolJson, String knownNamesJson, String toolCallsJson) throws Exception { String script = JavaScriptBuilder.enrichToolsScript( - httpJson, "{}", "{}", agentToolJson, "{}", "{}", "{}", "{}", knownNamesJson, "{}"); + httpJson, "{}", "{}", agentToolJson, "{}", "{}", "{}", "{}", knownNamesJson); // Wrap so the script's IIFE return is captured AND we get a JSON string // back — Graal's Value.toString() is JS source, not JSON. String wrapped = "var $ = {" diff --git a/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/compiler/AgentCompiler.java b/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/compiler/AgentCompiler.java index 513a364b..6fe6d2b7 100644 --- a/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/compiler/AgentCompiler.java +++ b/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/compiler/AgentCompiler.java @@ -359,12 +359,30 @@ public static Map> collectToolCredentials(AgentConfig confi } } List effective = own.isEmpty() ? agentCreds : own; - if (!effective.isEmpty()) map.put(tool.getName(), new ArrayList<>(effective)); + if (!effective.isEmpty()) map.put(tool.getName(), new ArrayList<>(new LinkedHashSet<>(effective))); } } return map; } + /** + * Collect the agent-level credential names, deduped and order-preserving. Used by + * {@code AgentService} to declare {@code TaskDef.runtimeMetadata} (embedded) on the non-worker + * SIMPLE tasks that run user-authored code — guardrails, callbacks, stop_when, gates, instructions, + * routers, graph node/edge workers — none of which carry their own per-item credential list, so the + * agent-level list is their only source. The host resolves the names at each task's poll and injects + * the values onto the wire-only {@code Task.runtimeMetadata}. + * + *

Note: the SDK worker wrappers for these non-worker task kinds do not yet read + * {@code Task.runtimeMetadata} (only the tool worker does), so declaring it here is currently inert — + * the values ride the wire but {@code get_secret()} inside those user functions won't resolve until + * the SDK wrappers are taught to route {@code runtimeMetadata} into the credential context.

+ */ + public static List collectAgentCredentials(AgentConfig config) { + if (config.getCredentials() == null || config.getCredentials().isEmpty()) return List.of(); + return new ArrayList<>(new LinkedHashSet<>(config.getCredentials())); + } + WorkflowDef compileWithTools(AgentConfig config) { ParsedModel parsed = ModelParser.parse(config.getModel()); String llmRef = toRef(config.getName()) + "_llm"; diff --git a/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/service/AgentService.java b/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/service/AgentService.java index 421b3427..52c192bc 100644 --- a/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/service/AgentService.java +++ b/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/service/AgentService.java @@ -951,24 +951,29 @@ private void registerTaskDefinitions(AgentConfig config) { @SuppressWarnings("unchecked") private void collectAndRegisterTasks(AgentConfig config, Set registered) { + // Credential names declared on each SIMPLE task's TaskDef.runtimeMetadata (embedded only, gated + // inside registerTaskDef). Worker tools use their per-tool creds (with agent-level fallback); + // the other user-code task kinds (guardrail/callback/stop_when/gate/instructions/router/graph) + // have no per-item credential list, so they use the agent-level names. Hoisted once per config. + Map> toolCreds = AgentCompiler.collectToolCredentials(config); + List agentCreds = AgentCompiler.collectAgentCredentials(config); + // Register dispatch task for this agent's tools if (config.getTools() != null) { for (ToolConfig tool : config.getTools()) { String tt = tool.getToolType(); if ("worker".equals(tt) && !registered.contains(tool.getName())) { - registerTaskDef( - tool.getName(), - AgentCompiler.collectToolCredentials(config).get(tool.getName())); + registerTaskDef(tool.getName(), toolCreds.get(tool.getName())); registered.add(tool.getName()); } } } - // Register stop_when worker + // Register stop_when worker (user-authored predicate → agent-level creds) if (config.getStopWhen() != null && config.getStopWhen().getTaskName() != null) { String taskName = config.getStopWhen().getTaskName(); if (!registered.contains(taskName)) { - registerTaskDef(taskName); + registerTaskDef(taskName, agentCreds); registered.add(taskName); } } @@ -987,7 +992,7 @@ private void collectAndRegisterTasks(AgentConfig config, Set registered) for (GuardrailConfig g : config.getGuardrails()) { if ("custom".equals(g.getGuardrailType()) && g.getTaskName() != null) { if (!registered.contains(g.getTaskName())) { - registerTaskDef(g.getTaskName()); + registerTaskDef(g.getTaskName(), agentCreds); registered.add(g.getTaskName()); } } @@ -998,7 +1003,7 @@ private void collectAndRegisterTasks(AgentConfig config, Set registered) if (config.getCallbacks() != null) { for (CallbackConfig cb : config.getCallbacks()) { if (cb.getTaskName() != null && !registered.contains(cb.getTaskName())) { - registerTaskDef(cb.getTaskName()); + registerTaskDef(cb.getTaskName(), agentCreds); registered.add(cb.getTaskName()); } } @@ -1007,7 +1012,7 @@ private void collectAndRegisterTasks(AgentConfig config, Set registered) // Register callable gate workers (text_contains gates are INLINE, no registration needed) if (config.getGate() != null && config.getGate().get("taskName") instanceof String gateTaskName) { if (!registered.contains(gateTaskName)) { - registerTaskDef(gateTaskName); + registerTaskDef(gateTaskName, agentCreds); registered.add(gateTaskName); } } @@ -1017,7 +1022,7 @@ private void collectAndRegisterTasks(AgentConfig config, Set registered) && instrMap.get("_worker_ref") instanceof String instrTaskName && !instrTaskName.isBlank()) { if (!registered.contains(instrTaskName)) { - registerTaskDef(instrTaskName); + registerTaskDef(instrTaskName, agentCreds); registered.add(instrTaskName); } } @@ -1026,12 +1031,12 @@ private void collectAndRegisterTasks(AgentConfig config, Set registered) if (config.getRouter() instanceof Map routerMap && routerMap.get("taskName") instanceof String routerTaskName) { if (!registered.contains(routerTaskName)) { - registerTaskDef(routerTaskName); + registerTaskDef(routerTaskName, agentCreds); registered.add(routerTaskName); } } else if (config.getRouter() instanceof WorkerRef workerRef && workerRef.getTaskName() != null) { if (!registered.contains(workerRef.getTaskName())) { - registerTaskDef(workerRef.getTaskName()); + registerTaskDef(workerRef.getTaskName(), agentCreds); registered.add(workerRef.getTaskName()); } } @@ -1120,7 +1125,7 @@ private void collectAndRegisterTasks(AgentConfig config, Set registered) for (Object nodeObj : nodes) { if (nodeObj instanceof Map node && node.get("_worker_ref") instanceof String workerRef) { if (!registered.contains(workerRef)) { - registerTaskDef(workerRef); + registerTaskDef(workerRef, agentCreds); registered.add(workerRef); } } @@ -1131,7 +1136,7 @@ private void collectAndRegisterTasks(AgentConfig config, Set registered) for (Object ceObj : condEdges) { if (ceObj instanceof Map ce && ce.get("_router_ref") instanceof String routerRef) { if (!registered.contains(routerRef)) { - registerTaskDef(routerRef); + registerTaskDef(routerRef, agentCreds); registered.add(routerRef); } } diff --git a/server/conductor-agentspan/src/test/java/dev/agentspan/runtime/compiler/WorkerRuntimeMetadataTest.java b/server/conductor-agentspan/src/test/java/dev/agentspan/runtime/compiler/WorkerRuntimeMetadataTest.java index 267ff6a6..380a371d 100644 --- a/server/conductor-agentspan/src/test/java/dev/agentspan/runtime/compiler/WorkerRuntimeMetadataTest.java +++ b/server/conductor-agentspan/src/test/java/dev/agentspan/runtime/compiler/WorkerRuntimeMetadataTest.java @@ -5,12 +5,14 @@ package dev.agentspan.runtime.compiler; import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.Mockito.atLeastOnce; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; import java.lang.reflect.Field; import java.lang.reflect.Method; +import java.util.HashMap; import java.util.List; import java.util.Map; @@ -24,6 +26,8 @@ import com.netflix.conductor.service.MetadataService; import dev.agentspan.runtime.model.AgentConfig; +import dev.agentspan.runtime.model.GuardrailConfig; +import dev.agentspan.runtime.model.TerminationConfig; import dev.agentspan.runtime.model.ToolConfig; import dev.agentspan.runtime.service.AgentService; import dev.agentspan.runtime.util.EmbeddedMode; @@ -102,6 +106,109 @@ void collectToolCredentials_mapsWorkerToItsSecretNames() { assertThat(creds.get("gh")).containsExactlyInAnyOrder("GITHUB_TOKEN", "GH_APP_ID"); } + // ── Agent-level creds feed the non-worker user-code task defs (guardrail/callback/etc.) ── + + @Test + void collectAgentCredentials_returnsDedupedOrdered() { + AgentConfig config = AgentConfig.builder() + .name("a") + .model("openai/gpt-4o") + .credentials(List.of("A", "B", "A")) + .build(); + assertThat(AgentCompiler.collectAgentCredentials(config)).containsExactly("A", "B"); + } + + @Test + void collectAgentCredentials_emptyWhenNoneDeclared() { + AgentConfig config = + AgentConfig.builder().name("a").model("openai/gpt-4o").build(); + assertThat(AgentCompiler.collectAgentCredentials(config)).isEmpty(); + } + + /** + * Wiring test: embedded, {@code collectAndRegisterTasks} must declare the agent-level creds on a + * custom-guardrail worker's {@link TaskDef} (user code → needs secrets), but leave the declarative + * {@code _termination} def empty (no user function runs there). Fails until agent-level creds are + * threaded into the guardrail registration site. + */ + @Test + void embedded_declaresAgentCredsOnGuardrailButNotTermination() throws Exception { + new EmbeddedMode().setEmbedded(true); + AgentConfig config = AgentConfig.builder() + .name("a") + .model("openai/gpt-4o") + .credentials(List.of("DEMO_SECRET")) + .guardrails(List.of(GuardrailConfig.builder() + .guardrailType("custom") + .taskName("a_guard") + .build())) + .termination(TerminationConfig.builder().build()) + .build(); + + Map defs = registerAllTaskDefs(config); + + assertThat(defs.get("a_guard").getRuntimeMetadata()).containsExactly("DEMO_SECRET"); + assertThat(defs.get("a_termination").getRuntimeMetadata()).isNullOrEmpty(); + } + + @Test + void standalone_leavesNonWorkerRuntimeMetadataEmpty() throws Exception { + new EmbeddedMode().setEmbedded(false); + AgentConfig config = AgentConfig.builder() + .name("a") + .model("openai/gpt-4o") + .credentials(List.of("DEMO_SECRET")) + .guardrails(List.of(GuardrailConfig.builder() + .guardrailType("custom") + .taskName("a_guard") + .build())) + .build(); + + Map defs = registerAllTaskDefs(config); + + assertThat(defs.get("a_guard").getRuntimeMetadata()).isNullOrEmpty(); + } + + /** + * Drive {@link AgentService}'s private {@code registerTaskDefinitions(AgentConfig)} and return every + * {@link TaskDef} handed to {@code MetadataService.registerTaskDef}, keyed by task name — so a test + * can assert per-task-kind {@code runtimeMetadata}. + */ + private static Map registerAllTaskDefs(AgentConfig config) throws Exception { + MetadataDAO metadataDAO = mock(MetadataDAO.class); + MetadataService metadataService = mock(MetadataService.class); + + AgentService service = new AgentService( + mock(dev.agentspan.runtime.compiler.AgentCompiler.class), + mock(dev.agentspan.runtime.normalizer.NormalizerRegistry.class), + mock(com.netflix.conductor.dao.ExecutionDAO.class), + metadataDAO, + mock(com.netflix.conductor.core.execution.WorkflowExecutor.class), + mock(com.netflix.conductor.service.WorkflowService.class), + mock(dev.agentspan.runtime.service.AgentStreamRegistry.class), + mock(com.netflix.conductor.service.ExecutionService.class), + mock(dev.agentspan.runtime.util.ProviderValidator.class)); + + Field msField = AgentService.class.getDeclaredField("metadataService"); + msField.setAccessible(true); + msField.set(service, metadataService); + + Method m = AgentService.class.getDeclaredMethod("registerTaskDefinitions", AgentConfig.class); + m.setAccessible(true); + m.invoke(service, config); + + @SuppressWarnings("unchecked") + ArgumentCaptor> captor = ArgumentCaptor.forClass(List.class); + verify(metadataService, atLeastOnce()).registerTaskDef(captor.capture()); + Map byName = new HashMap<>(); + for (List batch : captor.getAllValues()) { + for (TaskDef def : batch) { + byName.put(def.getName(), def); + } + } + return byName; + } + /** * Drive {@link AgentService}'s private {@code registerTaskDef(String, List)} with the credential * names {@link AgentCompiler#collectToolCredentials} yields for {@code toolName}, and capture the From e9b439156c3b182adf73b7995d45e611bd4c08b2 Mon Sep 17 00:00:00 2001 From: Viren Baraiya Date: Fri, 10 Jul 2026 08:41:36 -0700 Subject: [PATCH 5/5] add runtime metadata --- .../credentials/AgentspanSecretsDAO.java | 92 +++++++++++++ .../CredentialDataSourceConfig.java | 2 +- .../credentials/CredentialSchemaMigrator.java | 2 +- .../EncryptedDbCredentialStoreProvider.java | 2 +- .../runtime/credentials/MasterKeyConfig.java | 2 +- .../src/main/resources/application.properties | 10 ++ .../credentials/AgentspanSecretsDAOTest.java | 122 ++++++++++++++++++ 7 files changed, 228 insertions(+), 4 deletions(-) create mode 100644 server/conductor-agentspan-server/src/main/java/dev/agentspan/runtime/credentials/AgentspanSecretsDAO.java create mode 100644 server/conductor-agentspan-server/src/test/java/dev/agentspan/runtime/credentials/AgentspanSecretsDAOTest.java diff --git a/server/conductor-agentspan-server/src/main/java/dev/agentspan/runtime/credentials/AgentspanSecretsDAO.java b/server/conductor-agentspan-server/src/main/java/dev/agentspan/runtime/credentials/AgentspanSecretsDAO.java new file mode 100644 index 00000000..af6023c7 --- /dev/null +++ b/server/conductor-agentspan-server/src/main/java/dev/agentspan/runtime/credentials/AgentspanSecretsDAO.java @@ -0,0 +1,92 @@ +/* + * Copyright (c) 2025 AgentSpan + * Licensed under the MIT License. + */ +package dev.agentspan.runtime.credentials; + +import java.util.List; +import java.util.stream.Collectors; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.stereotype.Component; + +import com.netflix.conductor.dao.SecretsDAO; + +import dev.agentspan.runtime.model.credentials.CredentialMeta; +import dev.agentspan.runtime.spi.CredentialStoreProvider; + +/** + * Bridges conductor's global {@link SecretsDAO} to AgentSpan's own {@link CredentialStoreProvider} + * (the encrypted credential store), scoped to the anonymous/system user. + * + *

Active only when {@code conductor.secrets.type=agentspan} — the "agentspan-as-host" mode where + * the AgentSpan server embeds conductor ({@code agentspan.embedded=true}) and also serves as + * the secret-resolving host. In that mode the embedded conductor's {@code RuntimeMetadataResolver} + * calls {@link #getSecret(String)} at each SIMPLE task's poll to resolve the secret names a worker + * declared on {@code TaskDef.runtimeMetadata}, injecting the resolved values onto the wire-only + * {@code Task.runtimeMetadata}. Selecting this DAO ({@code havingValue="agentspan"}) gates conductor's + * own env-variable / noop {@code SecretsDAO} implementations off (they require + * {@code conductor.secrets.type} to be {@code env}/absent or {@code noop}).

+ * + *

Conductor secrets are global (name only); AgentSpan's store is per-user, so every lookup is + * scoped to {@link #ANONYMOUS_USER_ID} — the no-auth/system user, matching {@code CredentialEnvSeeder} + * and {@code AuthFilter.ANONYMOUS}. Names are treated as flat keys (no dotted JSONPath): worker + * credential names are simple identifiers, and {@link CredentialStoreProvider#get} resolves them + * directly.

+ * + *

The backing store beans ({@code EncryptedDbCredentialStoreProvider}, {@code MasterKeyConfig}, + * {@code CredentialDataSourceConfig}, {@code CredentialSchemaMigrator}) are normally dormant when + * embedded; they are re-enabled under this same {@code conductor.secrets.type=agentspan} flag so this + * bridge has a store to read from.

+ */ +@Component +@ConditionalOnProperty(name = "conductor.secrets.type", havingValue = "agentspan") +public class AgentspanSecretsDAO implements SecretsDAO { + + private static final Logger log = LoggerFactory.getLogger(AgentspanSecretsDAO.class); + + /** + * User ID for the anonymous/OSS user — matches {@code CredentialEnvSeeder.ANONYMOUS_USER_ID} and + * {@code AuthFilter.ANONYMOUS}. Conductor's global secret lookups resolve against this user. + */ + static final String ANONYMOUS_USER_ID = "00000000-0000-0000-0000-000000000000"; + + private final CredentialStoreProvider store; + + public AgentspanSecretsDAO(CredentialStoreProvider store) { + this.store = store; + log.info( + "AgentspanSecretsDAO active — embedded conductor secrets resolve from the AgentSpan " + + "credential store (scoped to system user {})", + ANONYMOUS_USER_ID); + } + + @Override + public String getSecret(String key) { + return store.get(ANONYMOUS_USER_ID, key); + } + + @Override + public boolean secretExists(String key) { + return store.get(ANONYMOUS_USER_ID, key) != null; + } + + @Override + public List listSecretNames() { + return store.list(ANONYMOUS_USER_ID).stream() + .map(CredentialMeta::getName) + .collect(Collectors.toList()); + } + + @Override + public void putSecret(String key, String value) { + store.set(ANONYMOUS_USER_ID, key, value); + } + + @Override + public void deleteSecret(String key) { + store.delete(ANONYMOUS_USER_ID, key); + } +} diff --git a/server/conductor-agentspan-server/src/main/java/dev/agentspan/runtime/credentials/CredentialDataSourceConfig.java b/server/conductor-agentspan-server/src/main/java/dev/agentspan/runtime/credentials/CredentialDataSourceConfig.java index a9d3c2c9..dc7ef9db 100644 --- a/server/conductor-agentspan-server/src/main/java/dev/agentspan/runtime/credentials/CredentialDataSourceConfig.java +++ b/server/conductor-agentspan-server/src/main/java/dev/agentspan/runtime/credentials/CredentialDataSourceConfig.java @@ -51,7 +51,7 @@ *

PostgreSQL: uses {@code org.postgresql.Driver} with a larger pool (default 8).

*/ @Configuration -@ConditionalOnProperty(name = "agentspan.embedded", havingValue = "false", matchIfMissing = true) +@ConditionalOnProperty(name = "conductor.secrets.type", havingValue = "agentspan") public class CredentialDataSourceConfig { private static final Logger log = LoggerFactory.getLogger(CredentialDataSourceConfig.class); diff --git a/server/conductor-agentspan-server/src/main/java/dev/agentspan/runtime/credentials/CredentialSchemaMigrator.java b/server/conductor-agentspan-server/src/main/java/dev/agentspan/runtime/credentials/CredentialSchemaMigrator.java index 17011b1b..83c83c8c 100644 --- a/server/conductor-agentspan-server/src/main/java/dev/agentspan/runtime/credentials/CredentialSchemaMigrator.java +++ b/server/conductor-agentspan-server/src/main/java/dev/agentspan/runtime/credentials/CredentialSchemaMigrator.java @@ -31,7 +31,7 @@ * pre-release development builds.

*/ @Component -@ConditionalOnProperty(name = "agentspan.embedded", havingValue = "false", matchIfMissing = true) +@ConditionalOnProperty(name = "conductor.secrets.type", havingValue = "agentspan") public class CredentialSchemaMigrator { private static final Logger log = LoggerFactory.getLogger(CredentialSchemaMigrator.class); diff --git a/server/conductor-agentspan-server/src/main/java/dev/agentspan/runtime/credentials/EncryptedDbCredentialStoreProvider.java b/server/conductor-agentspan-server/src/main/java/dev/agentspan/runtime/credentials/EncryptedDbCredentialStoreProvider.java index 6bc8a38e..4fcd66fb 100644 --- a/server/conductor-agentspan-server/src/main/java/dev/agentspan/runtime/credentials/EncryptedDbCredentialStoreProvider.java +++ b/server/conductor-agentspan-server/src/main/java/dev/agentspan/runtime/credentials/EncryptedDbCredentialStoreProvider.java @@ -35,7 +35,7 @@ *

The master key is the 32-byte key from {@code MasterKeyConfig#credentialMasterKey()}.

*/ @Component -@ConditionalOnProperty(name = "agentspan.embedded", havingValue = "false", matchIfMissing = true) +@ConditionalOnProperty(name = "conductor.secrets.type", havingValue = "agentspan") public class EncryptedDbCredentialStoreProvider implements CredentialStoreProvider { private static final Logger log = LoggerFactory.getLogger(EncryptedDbCredentialStoreProvider.class); diff --git a/server/conductor-agentspan-server/src/main/java/dev/agentspan/runtime/credentials/MasterKeyConfig.java b/server/conductor-agentspan-server/src/main/java/dev/agentspan/runtime/credentials/MasterKeyConfig.java index 6a2bb937..ff1f7a3a 100644 --- a/server/conductor-agentspan-server/src/main/java/dev/agentspan/runtime/credentials/MasterKeyConfig.java +++ b/server/conductor-agentspan-server/src/main/java/dev/agentspan/runtime/credentials/MasterKeyConfig.java @@ -29,7 +29,7 @@ * */ @Configuration -@ConditionalOnProperty(name = "agentspan.embedded", havingValue = "false", matchIfMissing = true) +@ConditionalOnProperty(name = "conductor.secrets.type", havingValue = "agentspan") public class MasterKeyConfig { private static final Logger log = LoggerFactory.getLogger(MasterKeyConfig.class); diff --git a/server/conductor-agentspan-server/src/main/resources/application.properties b/server/conductor-agentspan-server/src/main/resources/application.properties index 3906a87f..ba92767f 100644 --- a/server/conductor-agentspan-server/src/main/resources/application.properties +++ b/server/conductor-agentspan-server/src/main/resources/application.properties @@ -162,6 +162,16 @@ agentspan.credentials.store=built-in agentspan.credentials.strict-mode=false agentspan.credentials.resolve.rate-limit=120 +# Secret backend for the embedded conductor (RuntimeMetadataResolver at task poll, and +# ${workflow.secrets.NAME} substitution). 'agentspan' backs it with AgentSpan's encrypted +# credential store via AgentspanSecretsDAO and activates the store beans (datasource, master +# key, schema migrator, store provider) — the same store the native credential services use. +# Defaulted on so the standalone server keeps its store; when embedded as the secret-resolving +# host, set agentspan.embedded=true and leave this at 'agentspan'. Override to conductor's own +# 'env'/'noop' backend only when the host delivers secrets and the native store is not wanted +# (the native credential services require the AgentSpan store, so do not override it standalone). +conductor.secrets.type=${CONDUCTOR_SECRETS_TYPE:agentspan} + # Mask secrets from the host-owned /api/workflow/{id} (raw Conductor) read path too. # Off by default so embedding this library never mutates a host's workflow responses; # AgentSpan's own /api/agent/* reads are always masked regardless of this flag. diff --git a/server/conductor-agentspan-server/src/test/java/dev/agentspan/runtime/credentials/AgentspanSecretsDAOTest.java b/server/conductor-agentspan-server/src/test/java/dev/agentspan/runtime/credentials/AgentspanSecretsDAOTest.java new file mode 100644 index 00000000..82ccbd72 --- /dev/null +++ b/server/conductor-agentspan-server/src/test/java/dev/agentspan/runtime/credentials/AgentspanSecretsDAOTest.java @@ -0,0 +1,122 @@ +/* + * Copyright (c) 2025 AgentSpan + * Licensed under the MIT License. + */ +package dev.agentspan.runtime.credentials; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.Mockito.mock; + +import java.util.ArrayList; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; + +import org.junit.jupiter.api.Test; +import org.springframework.boot.test.context.runner.ApplicationContextRunner; +import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.Import; + +import dev.agentspan.runtime.model.credentials.CredentialMeta; +import dev.agentspan.runtime.spi.CredentialStoreProvider; + +/** + * {@link AgentspanSecretsDAO} bridges conductor's global {@code SecretsDAO} to AgentSpan's per-user + * {@link CredentialStoreProvider}, scoped to the anonymous/system user. Verifies the name→value + * round-trip is scoped to {@code ANONYMOUS_USER_ID} (so other users' secrets are invisible) and that + * the bean is selected only by {@code conductor.secrets.type=agentspan}. + */ +class AgentspanSecretsDAOTest { + + private static final String ANON = "00000000-0000-0000-0000-000000000000"; + + /** In-memory {@link CredentialStoreProvider} keyed by (userId,name) so scope can be asserted. */ + static class FakeStore implements CredentialStoreProvider { + final Map data = new LinkedHashMap<>(); + + private static String k(String u, String n) { + return u + "|" + n; + } + + @Override + public String get(String userId, String name) { + return data.get(k(userId, name)); + } + + @Override + public void set(String userId, String name, String value) { + data.put(k(userId, name), value); + } + + @Override + public void delete(String userId, String name) { + data.remove(k(userId, name)); + } + + @Override + public List list(String userId) { + List out = new ArrayList<>(); + for (String key : data.keySet()) { + int bar = key.indexOf('|'); + if (key.substring(0, bar).equals(userId)) { + out.add(CredentialMeta.builder().name(key.substring(bar + 1)).build()); + } + } + return out; + } + } + + @Test + void roundTrip_scopedToAnonymousUser() { + FakeStore store = new FakeStore(); + AgentspanSecretsDAO dao = new AgentspanSecretsDAO(store); + + assertThat(dao.secretExists("GITHUB_TOKEN")).isFalse(); + assertThat(dao.getSecret("GITHUB_TOKEN")).isNull(); + + dao.putSecret("GITHUB_TOKEN", "ghp_x"); + // written under the anonymous/system user — the scope conductor resolves against + assertThat(store.data).containsEntry(ANON + "|GITHUB_TOKEN", "ghp_x"); + assertThat(dao.getSecret("GITHUB_TOKEN")).isEqualTo("ghp_x"); + assertThat(dao.secretExists("GITHUB_TOKEN")).isTrue(); + + dao.putSecret("SLACK_TOKEN", "xoxb"); + assertThat(dao.listSecretNames()).containsExactlyInAnyOrder("GITHUB_TOKEN", "SLACK_TOKEN"); + + dao.deleteSecret("GITHUB_TOKEN"); + assertThat(dao.getSecret("GITHUB_TOKEN")).isNull(); + assertThat(dao.listSecretNames()).containsExactly("SLACK_TOKEN"); + } + + @Test + void doesNotReadOtherUsersSecrets() { + FakeStore store = new FakeStore(); + store.set("some-other-user", "GITHUB_TOKEN", "not-mine"); + AgentspanSecretsDAO dao = new AgentspanSecretsDAO(store); + assertThat(dao.getSecret("GITHUB_TOKEN")).isNull(); + assertThat(dao.listSecretNames()).isEmpty(); + } + + // ── gating: selected only by conductor.secrets.type=agentspan ── + + @Configuration + @Import(AgentspanSecretsDAO.class) + static class DaoConfig {} + + private final ApplicationContextRunner runner = new ApplicationContextRunner() + .withBean(CredentialStoreProvider.class, () -> mock(CredentialStoreProvider.class)) + .withUserConfiguration(DaoConfig.class); + + @Test + void beanPresent_whenConductorSecretsTypeAgentspan() { + runner.withPropertyValues("conductor.secrets.type=agentspan") + .run(ctx -> assertThat(ctx).hasSingleBean(AgentspanSecretsDAO.class)); + } + + @Test + void beanAbsent_whenFlagUnsetOrDifferent() { + runner.run(ctx -> assertThat(ctx).doesNotHaveBean(AgentspanSecretsDAO.class)); + runner.withPropertyValues("conductor.secrets.type=env") + .run(ctx -> assertThat(ctx).doesNotHaveBean(AgentspanSecretsDAO.class)); + } +}