diff --git a/design/2026-07-09-embedded-secret-toggle-design.md b/design/2026-07-09-embedded-secret-toggle-design.md new file mode 100644 index 00000000..a1406264 --- /dev/null +++ b/design/2026-07-09-embedded-secret-toggle-design.md @@ -0,0 +1,257 @@ +# Secret delivery toggle: native (standalone) vs host-delivered (embedded) + +**Date:** 2026-07-09 · **Status:** In progress · **Branch:** `feature/embedded-secret-toggle` + +## Summary + +AgentSpan keeps its full native credential mechanism and toggles it with one flag, `agentspan.embedded`: + +| Deployment | `agentspan.embedded` | Secrets | +|---|---|---| +| **Standalone** agentspan server | `false` (default) | **Native**: encrypted store, execution-token minting, `/api/workers/secrets` pull, SDK fetchers. Unchanged from `main`. | +| **Embedded** in orkes-conductor / conductor-oss | `true` | **Native dormant** (beans gated off); the **host** resolves secrets. | + +Nothing is deleted — the native code stays intact for standalone. + +## How secrets are delivered when embedded (split by task type) + +- **Worker tools (SIMPLE, polled by the SDK)** → the worker's `TaskDef.runtimeMetadata` declares the + secret names; the host resolves them at poll and injects the values onto the **wire-only + `Task.runtimeMetadata`**. This is the **target**. Until the client SDKs expose that field (see + table), we ship an **interim**: the compiler stamps + `inputParameters.__resolved_credentials__ = {NAME: "${workflow.secrets.NAME}"}` and the host + resolves it — same delivery, but it rides in the task-input `Map` that today's clients already keep. +- **LLM provider keys → the host's AI integration** (not a workflow secret). The `LLM_CHAT_COMPLETE` + task resolves its model — and its API key — from the configured AI integration by provider name; + agentspan stamps nothing. See the sequence below. (We deliberately do **not** map a provider to a + conventionally-named workflow secret — that would duplicate and can conflict with the integration.) +- **HTTP / MCP / planner-context headers → `${workflow.secrets.NAME}`.** These are the *user's* + external-API secrets (not integration-managed). The compiler rewrites a `${NAME}` placeholder in a + header to `${workflow.secrets.NAME}` (embedded) and the host substitutes it in memory before the + in-process call. Same for target and interim. + +### LLM provider key — via the host AI integration + +Embedded in Orkes, `OrkesAIModelProvider` (`@Primary`) is the active `AIModelProvider`. It resolves +the model **per call** from the integration store, scoped to the org — the API key lives in the +integration config and the built model client, and never touches the workflow definition or task input. +Verified in `orkes-conductor` (`workers/.../integrations/OrkesAIModelProvider.java`, +`ModelConfigurationProvider.java`). + +```mermaid +sequenceDiagram + autonumber + participant OP as Operator / UI + participant IS as IntegrationService (Orkes store) + participant TK as LLM_CHAT_COMPLETE task + participant LW as LLMWorkers (host) + participant MP as OrkesAIModelProvider (@Primary) + participant MC as ModelConfigurationProvider + participant API as Provider API (OpenAI, ...) + + Note over OP,IS: setup - integration stored per org (api_key in its config) + OP->>IS: create AI integration (provider=openai, api_key=...) + + Note over TK,API: execution - per LLM call + TK->>LW: LLM_CHAT_COMPLETE (llmProvider, model, integrationNames[AI_MODEL]) + LW->>MP: getModel(input) + MP->>MP: orgId from taskId, integrationName from input.integrationNames[AI_MODEL] + MP->>IS: getIntegration(orgId, integrationName) + IS-->>MP: Integration.configuration (incl api_key) + MP->>MC: getConfiguration(type, configMap) - build AIModel with api_key (cached) + MC-->>MP: AIModel + MP-->>LW: AIModel + LW->>API: chatComplete(messages) using the integration key + API-->>LW: completion +``` + +Consequences: +- **Agentspan stamps nothing on the LLM task.** `OrkesAIModelProvider` never reads an `apiKey` from + task input — it resolves by `(orgId, integrationName)` — so the interim `injectCredentialReferences` + + `LlmProviderEnv` mapping was redundant *and* bypassed. Both are removed. +- **Standalone** (not embedded): agentspan's own `AgentspanAIModelProvider` resolves per-user keys from + the native store — a separate path, unchanged. +- **Conductor-OSS** (no integration store): the OSS `AIModelProvider` serves models from startup + `ModelConfiguration`s — still not from workflow secrets. + +## Interim worker path (enrichment) — how it actually works + +Worker tools aren't static: the LLM picks them, an **INLINE "enrich" task** (GraalJS) builds the +SIMPLE tasks at runtime, and a `FORK_JOIN_DYNAMIC` schedules them. The per-tool cred map is baked +into the enrich script at compile time. + +### Interim sequence (`__resolved_credentials__`) + +```mermaid +sequenceDiagram + autonumber + participant C as Compiler + participant WF as WorkflowDef + participant LLM as LLM task + participant EN as Enrich task (INLINE / GraalJS) + participant H as Host (Orkes secretsDAO) + participant FK as FORK_JOIN_DYNAMIC + participant W as SDK worker (SIMPLE task) + participant T as Tool fn + + Note over C,WF: compile / register time (embedded only) + C->>C: collectToolCredentials(agent) - tool creds, agent-level fallback + C->>C: buildWorkerCredConfig - map each tool to its workflow.secrets refs + C->>WF: bake the cred map into the enrich INLINE script + Note over LLM,T: execution time + LLM->>EN: toolCalls (which tools to run) + H->>EN: substituteSecrets resolves the secret refs to plaintext + Note right of EN: caveat - resolves here, not at the SIMPLE task poll + EN->>EN: set inputParameters.__resolved_credentials__ on each SIMPLE task + EN->>FK: dynamicTasks + FK->>W: schedule and poll the SIMPLE task + W->>W: read __resolved_credentials__, set CredentialContext, strip key + W->>T: run tool, then get_secret(NAME) returns the value + T-->>W: result +``` + +**Caveat:** because the reference is baked into the INLINE script, the host resolves +`${workflow.secrets.NAME}` **at the enrich step**, not at the SIMPLE task's poll. So plaintext lands +in the forked task's **persisted** input, and a secret with JS-special chars (`"`, `\`, newline) can +break the script. The target fixes both. + +## Target vs interim — same runtime cost, better safety + +`get_secret(NAME)` (and `getCredential` / `ToolContext.getCredential` / `Secrets.Get`) is an +**in-memory lookup**: the worker stashes the resolved `{NAME: value}` map in a per-invocation context +and the accessor reads it. When embedded, the value is delivered **inline with the task** (poll +response) in both paths — **no extra calls** (only the standalone native path calls +`/api/workers/secrets`). So the choice is about safety, not performance — the target wins on: + +1. **Wire-only, never persisted** — `Task.runtimeMetadata` is on the poll response only; the interim + bakes plaintext into the forked task's persisted input (visible in execution history). +2. **No JS-injection** — declared names on the TaskDef vs a `${...}` reference baked into GraalJS + (special chars break it). +3. **Resolved at the SIMPLE task's own poll**, scoped to that task — not early, in a shared enrich task. +4. **First-class & declarative** — the conductor-native field vs a magic `__resolved_credentials__` key. + +Cost: the target needs the client libraries to expose `Task.runtimeMetadata` first — which is why the +interim ships now. + +## Required change in each Conductor client SDK (blocks the target) + +`Task.runtimeMetadata` is a new top-level field; today's clients drop it on the wire (no field, no +catch-all deserializer). Add it to each client's `Task` model (JSON key `runtimeMetadata`, +string→string, output-empty omitted), release, and bump the SDK's client dependency. + +| AgentSpan SDK | Client dependency | Client repo | Change to the `Task` model | +|---|---|---|---| +| Java | `org.conductoross:conductor-client:5.0.1` | `conductor-oss/java-sdk` | Add `Map runtimeMetadata` + getter/setter to `com.netflix.conductor.common.metadata.tasks.Task` (Jackson auto-maps; `@JsonInclude(NON_EMPTY)`). | +| Python | `conductor-python>=1.3.11` | `conductor-oss/python-sdk` | In `.../models/task.py`: add `runtime_metadata` to `swagger_types`/`attribute_map` (`'runtimeMetadata'`) + property. | +| C# | `conductor-csharp:1.1.4` | `conductor-oss/csharp-sdk` | Add `Dictionary RuntimeMetadata` with `[DataMember(Name="runtimeMetadata", EmitDefaultValue=false)]`. | +| TypeScript | `@io-orkes/conductor-javascript:^3.0.3` | Orkes TS SDK | Add `runtimeMetadata?: Record` to the `Task` type (JS keeps unknown keys; this is a type-def change so the read-path compiles). | + +(Go is out of scope — AgentSpan ships no Go SDK.) + +## Implementation (done + tested; CI green) + +- **Gating:** `@ConditionalOnProperty(agentspan.embedded=false, matchIfMissing=true)` on every native + secret bean — `WorkerController`, `CredentialResolutionService`, `ExecutionTokenService`, + `CredentialAwareMcpService`, `CredentialMaskingResponseAdvice`, `EncryptedDbCredentialStoreProvider`, + `MasterKeyConfig`, `CredentialEnvSeeder`, `CredentialSchemaMigrator`, `CredentialDataSourceConfig`, + `NoOpSecretOutputMasker`. Active consumers made tolerant: `AgentspanAIModelProvider` (`ObjectProvider` + + guards); `AgentService` / `AgentEventListener` (`@Autowired(required=false)` + null guards, so token + minting is skipped). +- **System tasks:** LLM keys come from the host AI integration (`OrkesAIModelProvider`) — agentspan + stamps nothing (the old `injectCredentialReferences` + `LlmProviderEnv` were removed). HTTP/MCP/ + planner headers emit `${workflow.secrets.NAME}` via `ToolCompiler.rewriteCredentialPlaceholders`. +- **Worker tools (interim):** `ToolCompiler.buildWorkerCredConfig` + `JavaScriptBuilder` enrich + injection + `AgentCompiler.collectToolCredentials`. +- **SDK read-path (all 4):** prefer the host map, else native token-pull; feed the existing accessor; + strip the key. Interim reads `inputData.__resolved_credentials__`; target reads `task.runtimeMetadata`. +- **Tests (all fail-first validated):** `NativeSecretGatingTest`, `ToolCompilerWorkerCredTest` + (GraalJS-runs the enrich script), `ReadResolvedCredentialsTest` (Java), `test_resolved_credentials.py` + (Python), TS `credentials.test.ts`. + +## The target change — main files (once clients expose `runtimeMetadata`) + +Net effect: **declare** secret names on the TaskDef instead of **stamping** a value-reference into +task input; the enrich script stops touching credentials entirely, which deletes the +JS-injection/persistence caveat above. System tasks are untouched — LLM keys stay on the host AI +integration, HTTP/MCP/planner headers keep their `${workflow.secrets.NAME}` rewrite. + +### Target sequence (`TaskDef.runtimeMetadata`) + +```mermaid +sequenceDiagram + autonumber + participant C as Compiler + participant MD as Metadata (TaskDef) + participant LLM as LLM task + participant EN as Enrich task (INLINE) + participant FK as FORK_JOIN_DYNAMIC + participant H as Host (RuntimeMetadataResolver + secretsDAO) + participant W as SDK worker (SIMPLE task) + participant T as Tool fn + + Note over C,MD: compile / register time (embedded only) + C->>C: collectToolCredentials(agent) - tool creds, agent-level fallback + C->>MD: register worker TaskDef with runtimeMetadata = [NAMES] + Note over LLM,T: execution time + LLM->>EN: toolCalls (which tools to run) + EN->>FK: dynamicTasks (SIMPLE tasks, NO creds in input) + W->>H: poll SIMPLE task + H->>H: resolve TaskDef.runtimeMetadata names to values (secretsDAO/env) + H-->>W: task, values on wire-only Task.runtimeMetadata (never persisted) + W->>W: read task.runtimeMetadata, set CredentialContext + W->>T: run tool, then get_secret(NAME) returns the value + T-->>W: result +``` + +Versus the interim: the enrich task never touches credentials, resolution happens at the SIMPLE +task's **own poll** (not the enrich step), and the value arrives on the **wire-only** +`Task.runtimeMetadata` — so nothing is baked into the script and no plaintext lands in persisted input. + +**0. Client libraries (prereq)** — add `Task.runtimeMetadata` per the table, release, and bump the +client dep in `sdk/java/build.gradle`, `sdk/python/pyproject.toml`, `sdk/csharp/.../Conductor.AI.csproj`, +`sdk/typescript/package.json`. + +**1. Server — declare, and stop stamping** (all embedded-gated on `EmbeddedMode.isEmbedded()`): + +| File | Change | +|---|---| +| `service/AgentService.java` | **ADD.** In `registerTaskDef`, set `taskDef.setRuntimeMetadata(names)` for each worker tool, where `names = AgentCompiler.collectToolCredentials(config).get(tool)`. This is the whole target delivery on the server. | +| `compiler/ToolCompiler.java` | **REMOVE** `buildWorkerCredConfig()` + `setWorkerCreds` + the `workerCredJson` argument passed to `enrichToolsScript` / `enrichToolsScriptDynamic`. | +| `util/JavaScriptBuilder.java` | **REMOVE** the `workerCredJson` param and the `if (workerCredCfg[n]) t.inputParameters.__resolved_credentials__ = …` lines in both enrich scripts. | +| `compiler/AgentCompiler.java`, `MultiAgentCompiler.java` | **MOVE.** Drop the `tc.setWorkerCreds(...)` calls; `collectToolCredentials` now feeds `AgentService` instead of `ToolCompiler`. | +| LLM keys | **UNCHANGED** — already handled by the host AI integration (`OrkesAIModelProvider`); no agentspan code. HTTP/MCP/planner headers keep their `${workflow.secrets}` rewrite. | + +**2. SDK worker read-path — read the field instead of the input key** (native token-pull fallback +stays in all four): + +| File | Change | +|---|---| +| `sdk/java/.../internal/WorkerManager.java` | `task.getRuntimeMetadata()` instead of `inputData.get("__resolved_credentials__")`. | +| `sdk/python/.../runtime/_dispatch.py` | `task.runtime_metadata` instead of `task.input_data.pop("__resolved_credentials__")`. | +| `sdk/csharp/.../WorkerManager.cs` | `task.RuntimeMetadata` instead of the `__resolved_credentials__` dict; drop the input-strip. | +| `sdk/typescript/src/worker.ts` | `task.runtimeMetadata` instead of `inputData["__resolved_credentials__"]`. The `credentials.ts` accessor (reads the resolved map from the context) is unchanged. | + +**3. Cleanup** — once all SDKs are on the new clients, delete the interim `__resolved_credentials__` +stamping (server) and reads (SDKs), plus `ToolCompilerWorkerCredTest`'s enrich-script assertions. + +The compiler/enrich change is a **deletion**; the real new code is one line in `AgentService` +(`setRuntimeMetadata`) plus a one-line read swap per SDK. Everything else (gating, system tasks, +accessors, native fallback) is already in place. + +## Dependency + +AgentSpan uses no PR #1255 API, so it builds/tests against the published `conductor 3.32.0-rc.3`. +`${workflow.secrets.NAME}` (and, in the target, `TaskDef.runtimeMetadata`) are resolved **at runtime +by the embedded host** (`substituteSecrets` / `RuntimeMetadataResolver` / `SecretsDAO`, PR #1255) — the +host must include PR #1255; agentspan does not build against it. (An earlier local +`…-runtimemeta-LOCAL` pin was reverted: its conductor-side `SecretResource` shadowed agentspan's +`SecretController` `GET /api/secrets` in standalone tests.) + +## Status + +| Item | State | +|---|---| +| Native mechanism gated on `agentspan.embedded` | ✅ done + tested | +| System tasks — LLM via host AI integration; HTTP/MCP/planner headers via `${workflow.secrets}` | ✅ done | +| Worker tools — interim `__resolved_credentials__` (server + 4 SDKs) | ✅ done + tested (CI green) | +| Worker tools — target `TaskDef.runtimeMetadata` | ⏳ blocked on client-SDK field (table) | 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 1c4d194a..8d8a7aee 100644 --- a/sdk/csharp/src/Conductor.AI/WorkerManager.cs +++ b/sdk/csharp/src/Conductor.AI/WorkerManager.cs @@ -105,8 +105,11 @@ 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. - Dictionary resolvedCredentials = new(); - if (_credentialNames.Length > 0) + // 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( toolCtx?.ExecutionToken, _credentialNames, ct); @@ -217,6 +220,24 @@ or CredentialRateLimitException } } + /// + /// 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 ReadRuntimeMetadata(Task task) + { + var result = new Dictionary(); + if (task?.RuntimeMetadata is { Count: > 0 } rm) + { + foreach (var (k, v) in rm) + if (k is not null && v is not null) + result[k] = v; + } + return result; + } + // ── JSON bridges (Newtonsoft ↔ System.Text.Json) ────────── /// Convert conductor-csharp's Newtonsoft-deserialized inputData to STJ JsonElements. 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 4c844443..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); @@ -373,7 +389,13 @@ 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()); - if (!declared.isEmpty()) { + // 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()) { String execToken = extractExecutionToken(inputData); try { resolvedSecrets = credentialFetcher.fetch(execToken, declared); @@ -417,6 +439,25 @@ private TaskResult executeHandler(String taskName, Task task) { return result; } + /** + * 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 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 : rm.entrySet()) { + if (e.getKey() != null && e.getValue() != null) { + out.put(e.getKey(), e.getValue()); + } + } + return out; + } + /** * Pull the execution token out of {@code inputData["__agentspan_ctx__"]["execution_token"]}. * Returns {@code null} if no token is present. 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/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/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/examples/demo_secret_resolution.py b/sdk/python/examples/demo_secret_resolution.py new file mode 100644 index 00000000..751e61d4 --- /dev/null +++ b/sdk/python/examples/demo_secret_resolution.py @@ -0,0 +1,118 @@ +# Copyright (c) 2025 Agentspan +# Licensed under the MIT License. See LICENSE file in the project root for details. + +"""Local e2e proof — embedded host (Orkes Conductor) resolves secrets. + +Minimal: ONE worker tool + an LLM call, so a single run exercises BOTH embedded +secret paths and the Orkes UI shows the result: + + * Worker-tool secret -> the tool declares credentials=["DEMO_SECRET"] and reads it + with get_secret(). Embedded, the compiler stamps + inputParameters.__resolved_credentials__ = {DEMO_SECRET: "${workflow.secrets.DEMO_SECRET}"} + and the host resolves it at poll. The tool returns a masked confirmation. + * LLM apiKey secret -> the agent makes an LLM call; embedded, the apiKey is stamped + ${workflow.secrets.} and resolved by the host. If the LLM step succeeds, + that secret resolved too. + +The tool does NO external network calls, so the task output is a clean, deterministic +"secret_resolved: true" you can screenshot. + +-------------------------------------------------------------------------------- +Setup (agentspan embedded in local Orkes, agentspan.embedded=true): + +1) In Orkes create the secrets (UI: Definitions -> Secrets, or the secrets API): + DEMO_SECRET = demo-value-12345 + OPENAI_API_KEY = # or the key matching AGENTSPAN_LLM_MODEL's provider + # (anthropic -> ANTHROPIC_API_KEY, etc.) + +2) Point the SDK at your local Orkes and give it an app key/secret: + export AGENTSPAN_SERVER_URL=http://localhost:8080/api + export AGENTSPAN_AUTH_KEY= + export AGENTSPAN_AUTH_SECRET= + export AGENTSPAN_LLM_MODEL=openai/gpt-4o + +3) Run: + cd sdk/python/examples + uv run python demo_secret_resolution.py + +Screenshot: in the Orkes UI open this execution -> the check_secret task's output +shows {"secret_resolved": true, "value_length": 16, "value_prefix": "demo…"}. + +CI smoke check: the script exits 0 only if the secret resolved AND the workflow +completed; otherwise it prints "SMOKE FAIL ..." and exits 1 (so it can gate CI). +""" + +from settings import settings + +from conductor.ai.agents import ( + Agent, + AgentRuntime, + CredentialNotFoundError, + get_secret, + tool, +) + + +@tool(credentials=["DEMO_SECRET"]) +def check_secret() -> dict: + """Report whether the declared secret was resolved by the host. No external calls. + + Returns a MASKED confirmation only — never logs or returns the full secret. + """ + try: + value = get_secret("DEMO_SECRET") + except CredentialNotFoundError: + return {"secret_resolved": False, "detail": "DEMO_SECRET was not delivered to the worker"} + return { + "secret_resolved": True, + "value_length": len(value), + "value_prefix": (value[:4] + "…") if value else "", + } + + +agent = Agent( + name="secret_resolution_demo", + model=settings.llm_model, + tools=[check_secret], + credentials=["DEMO_SECRET"], + instructions=( + "Call the check_secret tool exactly once, then state whether the secret " + "resolved and its length. Never invent or guess a secret value." + ), +) + + +def _secret_resolved(result) -> bool: + """True iff the check_secret tool ran and reported secret_resolved=True.""" + for call in result.tool_calls: + if call.get("name") != "check_secret" or "result" not in call: + continue + out = call["result"] + # Tool returns its dict directly; tolerate a {"result": {...}} wrapper too. + if isinstance(out, dict) and isinstance(out.get("result"), dict): + out = out["result"] + if isinstance(out, dict) and out.get("secret_resolved") is True: + return True + return False + + +if __name__ == "__main__": + import sys + + with AgentRuntime() as runtime: + result = runtime.run(agent, "Did my secret resolve? Use the tool to check.") + result.print_result() + + ok = result.is_success and _secret_resolved(result) + print("-" * 60) + if ok: + print("✅ SMOKE PASS — host resolved DEMO_SECRET (and the LLM call succeeded).") + sys.exit(0) + reason = ( + "workflow did not complete successfully" + if not result.is_success + else "check_secret did not report secret_resolved=true " + "(secret not delivered, tool not called, or running standalone)" + ) + print(f"❌ SMOKE FAIL — {reason}. status={result.status}") + sys.exit(1) 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 2fd569ed..a5188bb7 100644 --- a/sdk/python/src/conductor/ai/agents/runtime/_dispatch.py +++ b/sdk/python/src/conductor/ai/agents/runtime/_dispatch.py @@ -419,8 +419,16 @@ def tool_worker(task: Task) -> TaskResult: credential_names = list( _workflow_credentials.get(task.workflow_instance_id, []) ) + # 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 credential_names: + if isinstance(host_delivered, dict) and host_delivered: + resolved_secrets = { + k: v for k, v in host_delivered.items() if isinstance(v, str) + } + elif credential_names: token = _extract_execution_token(task) fetcher = _get_credential_fetcher() try: 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 diff --git a/sdk/python/tests/unit/test_runtime_metadata.py b/sdk/python/tests/unit/test_runtime_metadata.py new file mode 100644 index 00000000..6fe8e68c --- /dev/null +++ b/sdk/python/tests/unit/test_runtime_metadata.py @@ -0,0 +1,61 @@ +"""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 + +from conductor.ai.agents.runtime._dispatch import make_tool_worker +from conductor.ai.agents.runtime.credentials.accessor import get_secret +from conductor.ai.agents.tool import get_tool_def, tool +from conductor.client.http.models.task import Task + + +def _worker(): + @tool(credentials=["GITHUB_TOKEN"]) + def read_token() -> str: + return get_secret("GITHUB_TOKEN") + + td = get_tool_def(read_token) + return make_tool_worker(td.func, td.name, tool_def=td) + + +def test_prefers_host_delivered_runtime_metadata(): + wrapper = _worker() + task = Task() + task.input_data = {} + task.runtime_metadata = {"GITHUB_TOKEN": "ghp_host_resolved"} + task.workflow_instance_id = "wf" + task.task_id = "t" + + # The native fetcher must NOT be consulted when the host already delivered the map. + with patch("conductor.ai.agents.runtime._dispatch._get_credential_fetcher") as mock_fetcher: + result = wrapper(task) + + assert result.status == "COMPLETED" + assert result.output_data["result"] == "ghp_host_resolved" + mock_fetcher.assert_not_called() + + +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" + + class _Fetcher: + def fetch(self, token, names): + assert token == "tok" + assert names == ["GITHUB_TOKEN"] + return {"GITHUB_TOKEN": "ghp_native_pull"} + + with patch( + "conductor.ai.agents.runtime._dispatch._get_credential_fetcher", + return_value=_Fetcher(), + ): + result = wrapper(task) + + assert result.status == "COMPLETED" + assert result.output_data["result"] == "ghp_native_pull" diff --git a/sdk/typescript/src/credentials.ts b/sdk/typescript/src/credentials.ts index e9994309..2b974511 100644 --- a/sdk/typescript/src/credentials.ts +++ b/sdk/typescript/src/credentials.ts @@ -13,6 +13,10 @@ interface CredentialContext { serverUrl: string; headers: Record; executionToken: string; + // Pre-resolved name→value map. Embedded: the host resolves declared secrets at poll + // time and injects them onto task.runtimeMetadata; getCredential() reads them from here + // instead of pulling via the (dormant) execution-token endpoint. + resolved?: Record; } // AsyncLocalStorage scopes context per async-call chain so concurrent worker @@ -36,8 +40,9 @@ export function runWithCredentialContext( headers: Record, executionToken: string, fn: () => Promise, + resolved?: Record, ): Promise { - return _credentialStore.run({ serverUrl, headers, executionToken }, fn); + return _credentialStore.run({ serverUrl, headers, executionToken, resolved }, fn); } /** @@ -202,7 +207,17 @@ export async function getCredential(name: string): Promise { ); } + // Embedded / host-delivered: read from the pre-resolved map, no endpoint pull. + if (ctx.resolved && ctx.resolved[name] !== undefined) { + return ctx.resolved[name]; + } + const { serverUrl, headers, executionToken } = ctx; + // No token (embedded, native endpoint dormant) and not in the resolved map → the secret + // was not delivered. Surface as not-found (the intended off-host trim) rather than pulling. + if (!executionToken) { + throw new CredentialNotFoundError(name); + } const resolved = await resolveCredentials(serverUrl, headers, executionToken, [name]); const value = resolved[name]; diff --git a/sdk/typescript/src/worker.ts b/sdk/typescript/src/worker.ts index dd606468..d07e9a9a 100644 --- a/sdk/typescript/src/worker.ts +++ b/sdk/typescript/src/worker.ts @@ -1,4 +1,8 @@ -import { createConductorClient, TaskManager, NonRetryableException } from "@io-orkes/conductor-javascript"; +import { + createConductorClient, + TaskManager, + NonRetryableException, +} from "@io-orkes/conductor-javascript"; import type { ConductorWorker, Task, TaskResult } from "@io-orkes/conductor-javascript"; import type { ToolContext } from "./types.js"; import { TerminalToolError } from "./errors.js"; @@ -259,9 +263,16 @@ export class WorkerManager { * Queue a worker for the given task name. * Replaces any existing worker with the same task name. */ - addWorker(taskName: string, handler: WorkerHandler, credentials?: string[], domain?: string): void { + addWorker( + taskName: string, + handler: WorkerHandler, + credentials?: string[], + domain?: string, + ): void { // Track (taskName, domain) pairs — same name under different domains are distinct workers - const idx = this.pendingWorkers.findIndex((w) => w.taskName === taskName && w.domain === domain); + const idx = this.pendingWorkers.findIndex( + (w) => w.taskName === taskName && w.domain === domain, + ); if (idx >= 0) { this.pendingWorkers[idx] = { taskName, handler, credentials, domain }; } else { @@ -336,9 +347,7 @@ export class WorkerManager { leaseExtendEnabled: true, ...(pw.domain ? { domain: pw.domain } : {}), - async execute( - task: Task, - ): Promise> { + async execute(task: Task): Promise> { // Circuit breaker if (isCircuitBreakerOpen(pw.taskName)) { throw new NonRetryableException(`Circuit breaker open for ${pw.taskName}`); @@ -355,21 +364,20 @@ export class WorkerManager { cleaned["__workflowInstanceId__"] = task.workflowInstanceId; if (toolContext) cleaned["__toolContext__"] = toolContext; - // Credential setup + // Credential setup. Embedded: the worker's declared TaskDef.runtimeMetadata secret names are + // resolved by the host from its secret store at poll time and delivered on the wire-only + // Task.runtimeMetadata (never persisted). Prefer that map; otherwise fall back to the native + // execution-token pull (standalone). Resolution is up-front (no env mutation yet) — injection + // happens inside runHandler() via injectSecretsForInvocation so mutate-invoke-restore is + // atomic under a process lock. See docs/design/secret-injection-contract.md. const execToken = extractExecutionToken(inputData); + const hostDelivered = (task as { runtimeMetadata?: Record }) + .runtimeMetadata; - // Resolve credentials up-front (no env mutation yet). Injection happens - // inside runHandler() via injectSecretsForInvocation so the mutate- - // invoke-restore sequence is atomic under a process-wide lock. - // See docs/design/secret-injection-contract.md. let resolvedCredentials: Record = {}; - if (pw.credentials?.length) { - if (!execToken) { - throw new NonRetryableException( - `Required credentials not found: ${pw.credentials.join(", ")}. ` + - `No execution token available.`, - ); - } + if (hostDelivered && Object.keys(hostDelivered).length > 0) { + resolvedCredentials = hostDelivered; + } else if (pw.credentials?.length && execToken) { try { resolvedCredentials = await resolveCredentials( serverUrl, @@ -383,14 +391,13 @@ export class WorkerManager { ); } } + // else: no host delivery and no execution token — proceed with empty credentials; a + // tool that genuinely needs a secret fails via the accessor (the intended off-host trim). - const runHandler = async (): Promise< - Omit - > => { + const runHandler = async (): Promise> => { try { - let result = await injectSecretsForInvocation( - resolvedCredentials, - () => pw.handler(cleaned), + let result = await injectSecretsForInvocation(resolvedCredentials, () => + pw.handler(cleaned), ); // State mutation capture @@ -416,11 +423,17 @@ export class WorkerManager { } }; - // Scope credential context per-async-call so concurrent workers do not - // share (and clobber) module-level state. Runs even without an exec - // token so handlers see a consistent context shape. - if (execToken) { - return runWithCredentialContext(serverUrl, headers, execToken, runHandler); + // Scope credential context per-async-call so getCredential() sees the resolved + // (host-delivered or pulled) values and concurrent workers do not clobber each + // other's module-level state. + if (execToken || Object.keys(resolvedCredentials).length > 0) { + return runWithCredentialContext( + serverUrl, + headers, + execToken ?? "", + runHandler, + resolvedCredentials, + ); } return runHandler(); }, diff --git a/sdk/typescript/tests/unit/credentials.test.ts b/sdk/typescript/tests/unit/credentials.test.ts index 1f84852c..1a40d8bc 100644 --- a/sdk/typescript/tests/unit/credentials.test.ts +++ b/sdk/typescript/tests/unit/credentials.test.ts @@ -304,48 +304,83 @@ describe("runWithCredentialContext", () => { vi.restoreAllMocks(); }); - it.each([1, 2, 3])( - "isolates concurrent executions (run %i)", - async () => { - // Reproduce the worker race that breaks test_suite2_tool_calling: - // 1. Worker A enters context, starts handler. - // 2. Worker B enters context, finishes, exits. - // 3. Worker A's handler later calls getCredential — without per-async - // isolation, B's exit nulled A's context and getCredential throws. - // Test re-runs (1-3) to surface scheduling-dependent regressions. - vi.stubGlobal( - "fetch", - vi.fn().mockImplementation(async (_url, init: RequestInit) => { - const body = JSON.parse(String(init.body)); - // Echo the token back in the resolved value so we can verify isolation. - const result: Record = {}; - for (const n of body.names) result[n] = `${body.token}:${n}`; - return { ok: true, json: async () => result }; - }), - ); - - async function workerHandler(execToken: string, delayMs: number) { - return runWithCredentialContext(serverUrl, headers, execToken, async () => { - await new Promise((r) => setTimeout(r, delayMs)); - return getCredential("MY_KEY"); - }); - } - - const results = await Promise.all([ - workerHandler("tok-A", 30), - workerHandler("tok-B", 5), - workerHandler("tok-C", 20), - workerHandler("tok-D", 10), - workerHandler("tok-E", 15), - ]); - - expect(results).toEqual([ - "tok-A:MY_KEY", - "tok-B:MY_KEY", - "tok-C:MY_KEY", - "tok-D:MY_KEY", - "tok-E:MY_KEY", - ]); - }, - ); + it.each([1, 2, 3])("isolates concurrent executions (run %i)", async () => { + // Reproduce the worker race that breaks test_suite2_tool_calling: + // 1. Worker A enters context, starts handler. + // 2. Worker B enters context, finishes, exits. + // 3. Worker A's handler later calls getCredential — without per-async + // isolation, B's exit nulled A's context and getCredential throws. + // Test re-runs (1-3) to surface scheduling-dependent regressions. + vi.stubGlobal( + "fetch", + vi.fn().mockImplementation(async (_url, init: RequestInit) => { + const body = JSON.parse(String(init.body)); + // Echo the token back in the resolved value so we can verify isolation. + const result: Record = {}; + for (const n of body.names) result[n] = `${body.token}:${n}`; + return { ok: true, json: async () => result }; + }), + ); + + async function workerHandler(execToken: string, delayMs: number) { + return runWithCredentialContext(serverUrl, headers, execToken, async () => { + await new Promise((r) => setTimeout(r, delayMs)); + return getCredential("MY_KEY"); + }); + } + + const results = await Promise.all([ + workerHandler("tok-A", 30), + workerHandler("tok-B", 5), + workerHandler("tok-C", 20), + workerHandler("tok-D", 10), + workerHandler("tok-E", 15), + ]); + + expect(results).toEqual([ + "tok-A:MY_KEY", + "tok-B:MY_KEY", + "tok-C:MY_KEY", + "tok-D:MY_KEY", + "tok-E:MY_KEY", + ]); + }); +}); + +// ── host-delivered credentials (embedded: task.runtimeMetadata) ────────── + +describe("getCredential with host-delivered resolved map", () => { + const serverUrl = "https://api.test"; + const headers = {}; + + afterEach(() => { + clearCredentialContext(); + vi.restoreAllMocks(); + }); + + it("reads from the resolved map without pulling the endpoint", async () => { + const fetchSpy = vi.fn(); + vi.stubGlobal("fetch", fetchSpy); + // Embedded shape: no execution token, values pre-resolved by the host onto the context. + const value = await runWithCredentialContext( + serverUrl, + headers, + "", + async () => getCredential("GITHUB_TOKEN"), + { GITHUB_TOKEN: "ghp_host_resolved" }, + ); + expect(value).toBe("ghp_host_resolved"); + expect(fetchSpy).not.toHaveBeenCalled(); // native /workers/secrets pull is bypassed + }); + + it("throws NotFound for an undelivered secret with no token (off-host trim)", async () => { + const fetchSpy = vi.fn(); + vi.stubGlobal("fetch", fetchSpy); + await expect( + runWithCredentialContext(serverUrl, headers, "", async () => getCredential("MISSING"), { + GITHUB_TOKEN: "ghp_host_resolved", + }), + ).rejects.toBeInstanceOf(CredentialNotFoundError); + expect(fetchSpy).not.toHaveBeenCalled(); + }); }); 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); diff --git a/server/build.gradle b/server/build.gradle index be6f1de4..abfabb53 100644 --- a/server/build.gradle +++ b/server/build.gradle @@ -15,7 +15,11 @@ repositories { // ── Version catalog ────────────────────────────────────────────── ext { - 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.5' lombokVersion = '1.18.42' log4jVersion = '2.24.3' sqliteJdbcVersion = '3.47.0.0' 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 a57a61b8..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 @@ -9,6 +9,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Value; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Primary; @@ -50,6 +51,7 @@ *

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

*/ @Configuration +@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/CredentialEnvSeeder.java b/server/conductor-agentspan-server/src/main/java/dev/agentspan/runtime/credentials/CredentialEnvSeeder.java index 4a9bc395..45ed1887 100644 --- a/server/conductor-agentspan-server/src/main/java/dev/agentspan/runtime/credentials/CredentialEnvSeeder.java +++ b/server/conductor-agentspan-server/src/main/java/dev/agentspan/runtime/credentials/CredentialEnvSeeder.java @@ -15,6 +15,7 @@ import org.springframework.beans.factory.annotation.Value; import org.springframework.boot.ApplicationArguments; import org.springframework.boot.ApplicationRunner; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.stereotype.Component; import dev.agentspan.runtime.spi.CredentialStoreProvider; @@ -34,6 +35,7 @@ * (Vault, AWS SM, etc.) manage their own secrets.

*/ @Component +@ConditionalOnProperty(name = "agentspan.embedded", havingValue = "false", matchIfMissing = true) public class CredentialEnvSeeder implements ApplicationRunner { private static final Logger log = LoggerFactory.getLogger(CredentialEnvSeeder.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 cb4175f3..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 @@ -11,6 +11,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.boot.context.event.ApplicationReadyEvent; import org.springframework.context.event.EventListener; import org.springframework.jdbc.core.namedparam.NamedParameterJdbcTemplate; @@ -30,6 +31,7 @@ * pre-release development builds.

*/ @Component +@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 e9b1f4cd..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 @@ -18,6 +18,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.dao.EmptyResultDataAccessException; import org.springframework.jdbc.core.namedparam.NamedParameterJdbcTemplate; import org.springframework.stereotype.Component; @@ -34,6 +35,7 @@ *

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

*/ @Component +@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 3eff09e6..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 @@ -15,6 +15,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Value; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @@ -28,6 +29,7 @@ * */ @Configuration +@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/java/dev/agentspan/runtime/credentials/NoOpSecretOutputMasker.java b/server/conductor-agentspan-server/src/main/java/dev/agentspan/runtime/credentials/NoOpSecretOutputMasker.java index 30492724..e225df95 100644 --- a/server/conductor-agentspan-server/src/main/java/dev/agentspan/runtime/credentials/NoOpSecretOutputMasker.java +++ b/server/conductor-agentspan-server/src/main/java/dev/agentspan/runtime/credentials/NoOpSecretOutputMasker.java @@ -4,6 +4,7 @@ */ package dev.agentspan.runtime.credentials; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.stereotype.Service; import dev.agentspan.runtime.spi.SecretOutputMasker; @@ -20,6 +21,7 @@ * containing newlines, quotes, or other JSON-escaped characters are still caught). */ @Service +@ConditionalOnProperty(name = "agentspan.embedded", havingValue = "false", matchIfMissing = true) public class NoOpSecretOutputMasker implements SecretOutputMasker { @Override diff --git a/server/conductor-agentspan-server/src/main/resources/application.properties b/server/conductor-agentspan-server/src/main/resources/application.properties index 4ec2be1f..7f6729aa 100644 --- a/server/conductor-agentspan-server/src/main/resources/application.properties +++ b/server/conductor-agentspan-server/src/main/resources/application.properties @@ -154,10 +154,27 @@ agentspan.skills.max-file-count=${AGENTSPAN_SKILLS_MAX_FILE_COUNT:2000} # ============================================================================= # Credential Store Configuration # ============================================================================= +# Deployment mode toggle. false = standalone: AgentSpan's native credential +# mechanism is ACTIVE (encrypted store, execution-token minting, +# /api/workers/secrets pull, SDK fetchers). true = embedded in a host +# (orkes-conductor / conductor-oss): the native mechanism is DORMANT (all its +# beans are gated off) and the host delivers secrets — worker tools via +# TaskDef.runtimeMetadata, system tasks via ${workflow.secrets.NAME}. +agentspan.embedded=false 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)); + } +} diff --git a/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/ai/AgentChatCompleteTaskMapper.java b/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/ai/AgentChatCompleteTaskMapper.java index b2479295..bafb825f 100644 --- a/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/ai/AgentChatCompleteTaskMapper.java +++ b/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/ai/AgentChatCompleteTaskMapper.java @@ -111,8 +111,8 @@ protected TaskModel getMappedTask(TaskMapperContext taskMapperContext) throws Te TaskModel taskModel = super.getMappedTask(taskMapperContext); WorkflowModel workflowModel = taskMapperContext.getWorkflowModel(); - // Per-user LLM key resolution is handled by AgentspanAIModelProvider.getModel() - // which creates a fresh AIModel with the user's credential. No inputData injection needed. + // LLM credentials are the host's concern: embedded, the AI integration supplies the key + // (OrkesAIModelProvider); standalone, AgentspanAIModelProvider resolves it. Nothing to stamp here. try { ChatCompletion chatCompletion = objectMapper.convertValue(taskModel.getInputData(), ChatCompletion.class); diff --git a/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/ai/AgentspanAIModelProvider.java b/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/ai/AgentspanAIModelProvider.java index 492c7fff..18af1e35 100644 --- a/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/ai/AgentspanAIModelProvider.java +++ b/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/ai/AgentspanAIModelProvider.java @@ -24,6 +24,8 @@ import org.conductoross.conductor.ai.providers.perplexity.PerplexityAIConfiguration; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.ObjectProvider; +import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.annotation.Primary; import org.springframework.core.env.Environment; import org.springframework.stereotype.Component; @@ -85,6 +87,27 @@ public class AgentspanAIModelProvider extends AIModelProvider { private final CredentialResolutionService resolutionService; private final ExecutionTokenService tokenService; + /** + * Spring constructor. The native credential beans are optional: when + * {@code agentspan.embedded=true} they are gated off (the host delivers secrets), + * so they resolve to {@code null} here and native per-user resolution is skipped. + */ + @Autowired + public AgentspanAIModelProvider( + List> modelConfigurations, + Environment env, + OkHttpClient conductorAiHttpClient, + ObjectProvider resolutionService, + ObjectProvider tokenService) { + this( + modelConfigurations, + env, + conductorAiHttpClient, + resolutionService.getIfAvailable(), + tokenService.getIfAvailable()); + } + + /** Direct constructor (used by tests, and by the Spring constructor above). */ public AgentspanAIModelProvider( List> modelConfigurations, Environment env, @@ -95,7 +118,9 @@ public AgentspanAIModelProvider( this.conductorAiHttpClient = conductorAiHttpClient; this.resolutionService = resolutionService; this.tokenService = tokenService; - log.info("AgentspanAIModelProvider initialized (per-user credential resolution enabled)"); + log.info( + "AgentspanAIModelProvider initialized (native per-user credential resolution {})", + resolutionService != null ? "enabled" : "disabled — embedded/host-delivered"); } @Override @@ -148,6 +173,7 @@ public AIModel getModel(LLMWorkerInput input) { * @return per-user API key, or null if not found */ private String resolveUserApiKey(String provider) { + if (resolutionService == null) return null; // native resolution gated off (embedded) String envVarName = PROVIDER_TO_ENV_VAR.get(provider.toLowerCase()); if (envVarName == null) return null; @@ -177,6 +203,7 @@ private String resolveUserApiKey(String provider) { */ @SuppressWarnings("unchecked") private String extractUserIdFromTaskContext() { + if (tokenService == null) return null; // native token service gated off (embedded) try { TaskContext ctx = TaskContext.get(); if (ctx == null || ctx.getTask() == null) return null; @@ -224,6 +251,7 @@ public String resolveConfiguredBaseUrl(String provider) { * Resolve any named credential for the current user. */ private String resolveUserCredential(String credentialName) { + if (resolutionService == null) return null; // native resolution gated off (embedded) String userId = extractUserIdFromTaskContext(); if (userId == null) { userId = RequestContextHolder.get().map(ctx -> ctx.getUserId()).orElse(null); 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 8f1a9874..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 @@ -340,6 +340,49 @@ WorkflowDef compileSimple(AgentConfig config) { // ── Agent with tools ──────────────────────────────────────────── + /** + * Collect {@code toolName -> [credentialNames]} for the agent's tools: each tool's own declared + * 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}. + */ + public static Map> collectToolCredentials(AgentConfig config) { + List agentCreds = config.getCredentials() != null ? config.getCredentials() : List.of(); + Map> map = new LinkedHashMap<>(); + if (config.getTools() != null) { + for (ToolConfig tool : config.getTools()) { + if (tool.getName() == null) continue; + List own = new ArrayList<>(); + if (tool.getConfig() != null && tool.getConfig().get("credentials") instanceof List cl) { + for (Object c : cl) { + if (c instanceof String s) own.add(s); + } + } + List effective = own.isEmpty() ? agentCreds : own; + 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/controller/CredentialMaskingResponseAdvice.java b/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/controller/CredentialMaskingResponseAdvice.java index 6e22fbf6..4be523e4 100644 --- a/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/controller/CredentialMaskingResponseAdvice.java +++ b/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/controller/CredentialMaskingResponseAdvice.java @@ -10,6 +10,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Value; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.core.MethodParameter; import org.springframework.http.MediaType; import org.springframework.http.converter.HttpMessageConverter; @@ -54,6 +55,7 @@ * depth — it should never block a response from going out.

*/ @ControllerAdvice +@ConditionalOnProperty(name = "agentspan.embedded", havingValue = "false", matchIfMissing = true) public class CredentialMaskingResponseAdvice implements ResponseBodyAdvice { private static final Logger log = LoggerFactory.getLogger(CredentialMaskingResponseAdvice.class); diff --git a/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/controller/WorkerController.java b/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/controller/WorkerController.java index c0fb9f0b..f8813d0a 100644 --- a/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/controller/WorkerController.java +++ b/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/controller/WorkerController.java @@ -13,6 +13,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Value; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.http.ResponseEntity; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.web.bind.annotation.*; @@ -43,6 +44,7 @@ @RestController @RequestMapping("/api/workers") @RequiredArgsConstructor +@ConditionalOnProperty(name = "agentspan.embedded", havingValue = "false", matchIfMissing = true) public class WorkerController { private static final Logger log = LoggerFactory.getLogger(WorkerController.class); diff --git a/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/credentials/CredentialAwareMcpService.java b/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/credentials/CredentialAwareMcpService.java index bf99ebef..7a4fcf5b 100644 --- a/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/credentials/CredentialAwareMcpService.java +++ b/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/credentials/CredentialAwareMcpService.java @@ -13,6 +13,7 @@ import org.conductoross.conductor.ai.mcp.MCPService; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.context.annotation.Primary; import org.springframework.stereotype.Component; @@ -42,6 +43,7 @@ */ @Component @Primary +@ConditionalOnProperty(name = "agentspan.embedded", havingValue = "false", matchIfMissing = true) public class CredentialAwareMcpService extends MCPService { private static final Logger log = LoggerFactory.getLogger(CredentialAwareMcpService.class); diff --git a/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/credentials/CredentialResolutionService.java b/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/credentials/CredentialResolutionService.java index 909f4db6..7d83e526 100644 --- a/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/credentials/CredentialResolutionService.java +++ b/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/credentials/CredentialResolutionService.java @@ -6,6 +6,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.stereotype.Service; import com.fasterxml.jackson.databind.JsonNode; @@ -39,6 +40,7 @@ * own {@code os.environ} fallback when {@code secret_strict_mode=false}.

*/ @Service +@ConditionalOnProperty(name = "agentspan.embedded", havingValue = "false", matchIfMissing = true) public class CredentialResolutionService { private static final Logger log = LoggerFactory.getLogger(CredentialResolutionService.class); diff --git a/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/credentials/ExecutionTokenService.java b/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/credentials/ExecutionTokenService.java index ec778b5a..18d87e67 100644 --- a/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/credentials/ExecutionTokenService.java +++ b/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/credentials/ExecutionTokenService.java @@ -15,6 +15,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Service; @@ -31,6 +32,7 @@ * (bounded risk: tokens expire with workflow TTL).

*/ @Service +@ConditionalOnProperty(name = "agentspan.embedded", havingValue = "false", matchIfMissing = true) public class ExecutionTokenService { private static final Logger log = LoggerFactory.getLogger(ExecutionTokenService.class); diff --git a/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/service/AgentEventListener.java b/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/service/AgentEventListener.java index 96ab32c4..a094cf6c 100644 --- a/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/service/AgentEventListener.java +++ b/server/conductor-agentspan/src/main/java/dev/agentspan/runtime/service/AgentEventListener.java @@ -266,6 +266,7 @@ private void handleWorkflowTerminated(WorkflowModel workflow) { } private void revokeWorkflowToken(WorkflowModel workflow) { + if (executionTokenService == null) return; // native token service gated off (embedded) try { Object ctx = workflow.getVariables() != null ? workflow.getVariables().get("__agentspan_ctx__") : null; 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 d193e44c..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 @@ -48,6 +48,7 @@ import dev.agentspan.runtime.credentials.ExecutionTokenService; import dev.agentspan.runtime.model.*; import dev.agentspan.runtime.normalizer.NormalizerRegistry; +import dev.agentspan.runtime.util.EmbeddedMode; import dev.agentspan.runtime.util.ModelParser; import dev.agentspan.runtime.util.ProviderValidator; import dev.agentspan.runtime.util.WorkflowClassifiers; @@ -950,22 +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()); + 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); } } @@ -984,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()); } } @@ -995,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()); } } @@ -1004,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); } } @@ -1014,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); } } @@ -1023,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()); } } @@ -1117,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); } } @@ -1128,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); } } @@ -1409,6 +1417,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 +1435,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/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..380a371d --- /dev/null +++ b/server/conductor-agentspan/src/test/java/dev/agentspan/runtime/compiler/WorkerRuntimeMetadataTest.java @@ -0,0 +1,246 @@ +/* + * 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.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; + +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.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; + +/** + * 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"); + } + + // ── 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 + * {@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); + } +} diff --git a/server/conductor-agentspan/src/test/java/dev/agentspan/runtime/credentials/NativeSecretGatingTest.java b/server/conductor-agentspan/src/test/java/dev/agentspan/runtime/credentials/NativeSecretGatingTest.java new file mode 100644 index 00000000..1176fae3 --- /dev/null +++ b/server/conductor-agentspan/src/test/java/dev/agentspan/runtime/credentials/NativeSecretGatingTest.java @@ -0,0 +1,63 @@ +/* + * 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 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.controller.WorkerController; +import dev.agentspan.runtime.spi.CredentialStoreProvider; + +/** + * Verifies the native secret mechanism toggles on {@code agentspan.embedded}: + * its beans are present in standalone mode (flag absent or {@code false}) and + * gated OFF when embedded ({@code agentspan.embedded=true}), where the host + * delivers secrets instead. + */ +class NativeSecretGatingTest { + + /** + * Registers the gated native beans (their class-level {@code @ConditionalOnProperty} + * is evaluated on import) and supplies mock collaborators so they can be constructed + * when the condition allows. + */ + @Configuration + @Import({WorkerController.class, CredentialResolutionService.class, ExecutionTokenService.class}) + static class NativeBeans {} + + private final ApplicationContextRunner runner = new ApplicationContextRunner() + .withBean(CredentialStoreProvider.class, () -> mock(CredentialStoreProvider.class)) + .withBean("credentialMasterKey", byte[].class, () -> new byte[32]) + .withUserConfiguration(NativeBeans.class); + + @Test + void nativeBeans_present_whenFlagAbsent() { + runner.run(ctx -> assertThat(ctx) + .hasSingleBean(WorkerController.class) + .hasSingleBean(CredentialResolutionService.class) + .hasSingleBean(ExecutionTokenService.class)); + } + + @Test + void nativeBeans_present_whenStandalone() { + runner.withPropertyValues("agentspan.embedded=false").run(ctx -> assertThat(ctx) + .hasSingleBean(WorkerController.class) + .hasSingleBean(CredentialResolutionService.class) + .hasSingleBean(ExecutionTokenService.class)); + } + + @Test + void nativeBeans_dormant_whenEmbedded() { + runner.withPropertyValues("agentspan.embedded=true").run(ctx -> assertThat(ctx) + .doesNotHaveBean(WorkerController.class) + .doesNotHaveBean(CredentialResolutionService.class) + .doesNotHaveBean(ExecutionTokenService.class)); + } +}