diff --git a/Conductor.AI.Tests/AgentStatusMappingTests.cs b/Conductor.AI.Tests/AgentStatusMappingTests.cs index bfd8fc55..4e346bd5 100644 --- a/Conductor.AI.Tests/AgentStatusMappingTests.cs +++ b/Conductor.AI.Tests/AgentStatusMappingTests.cs @@ -26,30 +26,8 @@ namespace Conductor.AI.Tests; public sealed class AgentStatusMappingTests { - private sealed class StubHandler : HttpMessageHandler - { - private readonly Func _respond; - public StubHandler(Func respond) => _respond = respond; - - protected override Task SendAsync(HttpRequestMessage request, CancellationToken ct) - { - var (status, body) = _respond(request); - return Task.FromResult(new HttpResponseMessage(status) - { - Content = new StringContent(body, Encoding.UTF8, "application/json"), - }); - } - } - private static Configuration BuildConfig(Func respond) - { - var configuration = new Configuration { BasePath = "http://server/api" }; - configuration.ApiClient.RestClient = new RestClient(new RestClientOptions("http://server/api") - { - ConfigureMessageHandler = _ => new StubHandler(respond), - }); - return configuration; - } + => StubAgentServer.Configure(respond); private static (HttpStatusCode, string) RouteStatusAndExecution( HttpRequestMessage request, string statusBody, string executionBody = "{}") @@ -57,12 +35,7 @@ private static (HttpStatusCode, string) RouteStatusAndExecution( private static (HttpStatusCode, string) Route( HttpRequestMessage request, string statusBody, string executionBody, string workflowBody) - { - var path = request.RequestUri!.AbsolutePath; - if (path.EndsWith("/status")) return (HttpStatusCode.OK, statusBody); - if (path.Contains("/execution/")) return (HttpStatusCode.OK, executionBody); - return (HttpStatusCode.OK, workflowBody); - } + => StubAgentServer.Route(request, statusBody, executionBody, workflowBody); // ── AgentRuntime.GetStatusAsync ────────────────────────────────────── @@ -171,11 +144,12 @@ public async Task WaitAsync_ExtractsToolCallsFromWorkflowTasks() executionBody: "{}", workflowBody: """ {"tasks":[ - {"taskType":"echo","referenceTaskName":"call_echo_1", - "inputData":{"query":"hello","_internal":"drop me","ctx":"drop me too"}, + {"taskType":"echo","taskDefName":"echo","referenceTaskName":"call_ceOwlp7lQ_0__1", + "inputData":{"query":"hello","method":"echo","_agent_tool_name":"echo", + "_agent_state":{},"_internal":"drop me","ctx":"drop me too"}, "outputData":{"result":"echoed: hello"}}, - {"taskType":"LLM_CHAT_COMPLETE","referenceTaskName":"llm_1", - "inputData":{},"outputData":{"promptTokens":10}} + {"taskType":"LLM_CHAT_COMPLETE","taskDefName":"llm_chat_complete", + "referenceTaskName":"llm_1","inputData":{},"outputData":{"promptTokens":10}} ]} """)); diff --git a/Conductor.AI.Tests/AgentToolCallExtractionTests.cs b/Conductor.AI.Tests/AgentToolCallExtractionTests.cs new file mode 100644 index 00000000..81a0104b --- /dev/null +++ b/Conductor.AI.Tests/AgentToolCallExtractionTests.cs @@ -0,0 +1,378 @@ +/* + * Copyright 2024 Conductor Authors. + *

+ * Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + *

+ * http://www.apache.org/licenses/LICENSE-2.0 + *

+ * Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on + * an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the + * specific language governing permissions and limitations under the License. + */ +// AgentResult.ToolCalls named every tool after its Conductor task type, which is +// the tool's own name only for a worker tool, and detection keyed on a `call_` +// reference-name prefix, which is OpenAI's tool-call ID format. Fixtures carry the +// shape real payloads have: fork index, loop suffix and the dispatch markers. + +using System.Linq; +using System.Net; +using System.Net.Http; +using System.Text.Json; +using System.Text; +using System.Threading.Tasks; +using Conductor.Client; +using RestSharp; +using Xunit; + +namespace Conductor.AI.Tests; + +public sealed class AgentToolCallExtractionTests +{ + ///

Drive WaitAsync to completion over a stubbed server and return the built result. + private static async Task WaitWithTasksAsync( + string workflowBody, string statusBody = DefaultStatus) + { + var configuration = StubAgentServer.Configure( + request => StubAgentServer.Route(request, statusBody, "{}", workflowBody)); + var handle = new AgentHandle("e1", new OrkesAgentClient(configuration)); + return await handle.WaitAsync(); + } + + private const string DefaultStatus = """ + {"executionId":"e1","status":"COMPLETED","isComplete":true,"isRunning":false, + "isWaiting":false,"output":{"result":"done"}} + """; + + private static Dictionary ArgsOf(Dictionary toolCall) + => Assert.IsType>(toolCall["args"]); + + // ── Tool name comes from the payload, never from the task type ─────── + + [Fact] + public async Task WorkerTool_NameFromAgentToolName() + { + // The one kind the old taskType read got right; pinned so it stays right. + var result = await WaitWithTasksAsync(""" + {"tasks":[ + {"taskType":"get_weather","taskDefName":"get_weather", + "referenceTaskName":"call_PMnNIdOPvm9EQ8e6tn2kbxPY_0__1", + "inputData":{"_agent_tool_name":"get_weather","_agent_state":{}, + "method":"get_weather","city":"San Francisco"}, + "outputData":{"result":"Sunny, 72F"}} + ]} + """); + + var toolCall = Assert.Single(result.ToolCalls!); + Assert.Equal("get_weather", toolCall["name"]); + Assert.Equal(["city"], ArgsOf(toolCall).Keys); + Assert.Equal("Sunny, 72F", toolCall["result"]!.ToString()); + } + + [Theory] + // Every non-worker kind in ToolCompiler.TYPE_MAP: task type is the system type. + [InlineData("HTTP", "lookup_order")] + [InlineData("CALL_MCP_TOOL", "search_docs")] + [InlineData("SUB_WORKFLOW", "billing_agent")] + [InlineData("HUMAN", "escalate")] + [InlineData("GENERATE_IMAGE", "draw_chart")] + [InlineData("GENERATE_AUDIO", "read_aloud")] + [InlineData("GENERATE_VIDEO", "animate")] + [InlineData("GENERATE_PDF", "make_report")] + [InlineData("LLM_INDEX_TEXT", "index_docs")] + [InlineData("LLM_SEARCH_INDEX", "search_index")] + [InlineData("PULL_WORKFLOW_MESSAGES", "read_queue")] + public async Task SystemTaskTool_NameFromAgentToolName_NotTaskType(string taskType, string toolName) + { + var result = await WaitWithTasksAsync($$""" + {"tasks":[ + {"taskType":"{{taskType}}","taskDefName":"{{toolName}}", + "referenceTaskName":"{{toolName}}_0__1", + "inputData":{"_agent_tool_name":"{{toolName}}","query":"x"}, + "outputData":{"result":"ok"} + } + ]} + """); + + var toolCall = Assert.Single(result.ToolCalls!); + Assert.Equal(toolName, toolCall["name"]); + Assert.NotEqual(taskType, toolCall["name"]); + } + + [Fact] + public async Task NameFallsBackToMethod_WhenAgentToolNameAbsent() + { + // The dynamic-tools script sets no `_agent_tool_name`; MCP carries `method`. + var result = await WaitWithTasksAsync(""" + {"tasks":[ + {"taskType":"CALL_MCP_TOOL","taskDefName":"call_mcp_tool", + "referenceTaskName":"search_docs_0", + "inputData":{"mcpServer":"docs","method":"search_docs","arguments":{"q":"x"}}, + "outputData":{"result":["a"]}} + ]} + """); + + Assert.Equal("search_docs", Assert.Single(result.ToolCalls!)["name"]); + } + + [Fact] + public async Task NameFallsBackToTaskDefName_WhenNeitherMarkerPresent() + { + var result = await WaitWithTasksAsync(""" + {"tasks":[ + {"taskType":"GENERATE_IMAGE","taskDefName":"generate_image", + "referenceTaskName":"generate_image_0", + "inputData":{"prompt":"a cat"}, + "outputData":{"result":"http://img"}} + ]} + """); + + Assert.Equal("generate_image", Assert.Single(result.ToolCalls!)["name"]); + } + + // ── Detection no longer keys on the provider's tool-call ID ────────── + + [Fact] + public async Task AnthropicReferenceName_StillDetected() + { + // Anthropic IDs start `toolu_`, so the old `call_` prefix test found nothing. + var result = await WaitWithTasksAsync(""" + {"tasks":[ + {"taskType":"get_weather","taskDefName":"get_weather", + "referenceTaskName":"toolu_01A09q90qw90lq917835lq9_0__1", + "inputData":{"_agent_tool_name":"get_weather","_agent_state":{},"city":"Paris"}, + "outputData":{"result":"Rainy"}} + ]} + """); + + Assert.Equal("get_weather", Assert.Single(result.ToolCalls!)["name"]); + } + + [Fact] + public async Task WorkerToolWithoutToolNameMarker_DetectedByAgentState() + { + // Dynamic dispatch injects `_agent_state` only, and a worker's task type is + // its own name, so no allowlist of system types can catch it. + var result = await WaitWithTasksAsync(""" + {"tasks":[ + {"taskType":"echo","taskDefName":"echo","referenceTaskName":"toolu_xyz_0", + "inputData":{"_agent_state":{},"method":"echo","text":"hi"}, + "outputData":{"result":"hi"}} + ]} + """); + + Assert.Equal("echo", Assert.Single(result.ToolCalls!)["name"]); + } + + [Fact] + public async Task NonToolTasks_Excluded() + { + // Agent scaffolding, plus the INLINE task substituted for a hallucinated name. + var result = await WaitWithTasksAsync(""" + {"tasks":[ + {"taskType":"LLM_CHAT_COMPLETE","taskDefName":"llm_chat_complete", + "referenceTaskName":"llm_1","inputData":{},"outputData":{"promptTokens":10}}, + {"taskType":"FORK_JOIN_DYNAMIC","taskDefName":"fork","referenceTaskName":"fork_1", + "inputData":{},"outputData":{}}, + {"taskType":"JOIN","taskDefName":"join","referenceTaskName":"join_1", + "inputData":{},"outputData":{}}, + {"taskType":"INLINE","taskDefName":"made_up_tool","referenceTaskName":"made_up_tool", + "inputData":{"evaluatorType":"graaljs","expression":"...","errorMessage":"Unknown tool"}, + "outputData":{"result":"Unknown tool 'made_up_tool'.","is_error":true}}, + {"taskType":"SWITCH","taskDefName":"switch","referenceTaskName":"switch_1", + "inputData":{},"outputData":{}}, + {"taskType":"DO_WHILE","taskDefName":"loop","referenceTaskName":"loop_1", + "inputData":{},"outputData":{}}, + {"taskType":"SET_VARIABLE","taskDefName":"set_var","referenceTaskName":"set_var_1", + "inputData":{},"outputData":{}} + ]} + """); + + Assert.Null(result.ToolCalls); + } + + [Fact] + public async Task MultiAgentSetVariable_NotAToolCall() + { + // A coordinator's SET_VARIABLE carries `_agent_state`, so the marker alone is + // not enough: the worker case also needs taskType to equal taskDefName. + var result = await WaitWithTasksAsync(""" + {"tasks":[ + {"taskType":"SET_VARIABLE","taskDefName":"triage_init", + "referenceTaskName":"triage_init", + "inputData":{"conversation":"hello","_agent_state":{"turn":1}}, + "outputData":{}} + ]} + """); + + Assert.Null(result.ToolCalls); + } + + [Fact] + public async Task WorkerTaskWithoutDispatchMarker_NotAToolCall() + { + // A framework passthrough is a SIMPLE task the LLM never dispatched. + var result = await WaitWithTasksAsync(""" + {"tasks":[ + {"taskType":"claude_code","taskDefName":"claude_code", + "referenceTaskName":"_fw_task", + "inputData":{"prompt":"hi","session_id":"s1"}, + "outputData":{"result":"hello"}} + ]} + """); + + Assert.Null(result.ToolCalls); + } + + [Theory] + // The agent layer emits these for its own structure (sub-agent, strategy, router, + // plan approval). None is LLM-dispatched, so none carries a marker. + [InlineData("SUB_WORKFLOW", "billing_strategy")] + [InlineData("SUB_WORKFLOW", "triage_router")] + [InlineData("HUMAN", "plan_approval")] + public async Task AgentStructureReusingAToolTaskType_NotAToolCall(string taskType, string taskDefName) + { + var result = await WaitWithTasksAsync($$""" + {"tasks":[ + {"taskType":"{{taskType}}","taskDefName":"{{taskDefName}}", + "referenceTaskName":"0_{{taskDefName}}__1", + "inputData":{"prompt":"hello","media":[],"session_id":"s1", + "context":{"turn":1} + }, + "outputData":{"result":"handled"} + } + ]} + """); + + Assert.Null(result.ToolCalls); + Assert.Single(result.Events!); + } + + [Fact] + public async Task AgentAsTool_DetectedByDispatchMarker() + { + // Same type, but the dispatch script marked it, so it is a real tool call. + var result = await WaitWithTasksAsync(""" + {"tasks":[ + {"taskType":"SUB_WORKFLOW","taskDefName":"billing_agent_wf", + "referenceTaskName":"toolu_01xyz_0__1", + "inputData":{"_agent_tool_name":"billing_agent","prompt":"refund order A-1", + "session_id":"s1"}, + "outputData":{"result":"refunded"}} + ]} + """); + + Assert.Equal("billing_agent", Assert.Single(result.ToolCalls!)["name"]); + } + + [Fact] + public async Task InternalInputKeys_StrippedFromArgs() + { + var result = await WaitWithTasksAsync(""" + {"tasks":[ + {"taskType":"HTTP","taskDefName":"lookup","referenceTaskName":"lookup_0", + "inputData":{"_agent_tool_name":"lookup","_agent_state":{},"method":"lookup", + "ctx":"x","workerTag":"y","agentConfig":{},"order":"A-1"}, + "outputData":{"result":"shipped"}} + ]} + """); + + Assert.Equal(["order"], ArgsOf(Assert.Single(result.ToolCalls!)).Keys); + } + + // ── Events ────────────────────────────────────────────────────────── + + [Fact] + public async Task Events_ToolCallAndResultPerToolTask_ThenDone() + { + var result = await WaitWithTasksAsync(""" + {"tasks":[ + {"taskType":"HTTP","taskDefName":"lookup_order","referenceTaskName":"lookup_order_0", + "inputData":{"_agent_tool_name":"lookup_order","order":"A-1"}, + "outputData":{"result":"shipped"}} + ]} + """); + + Assert.Collection( + result.Events!, + e => + { + Assert.Equal(EventType.ToolCall, e.Type); + Assert.Equal("lookup_order", e.ToolName); + Assert.Equal(["order"], e.Args!.Keys); + }, + e => + { + Assert.Equal(EventType.ToolResult, e.Type); + Assert.Equal("lookup_order", e.ToolName); + Assert.Equal("shipped", e.Result!.ToString()); + }, + e => + { + Assert.Equal(EventType.Done, e.Type); + Assert.Equal("e1", e.ExecutionId); + }); + } + + [Fact] + public async Task Events_NeverNull_SoEnumerationDoesNotThrow() + { + var result = await WaitWithTasksAsync("""{"tasks":[]}"""); + + Assert.NotNull(result.Events); + var terminal = Assert.Single(result.Events!); + Assert.Equal(EventType.Done, terminal.Type); + } + + [Fact] + public async Task Events_FailedRun_EndsWithErrorNotDone() + { + var result = await WaitWithTasksAsync("""{"tasks":[]}""", statusBody: """ + {"executionId":"e1","status":"FAILED","isComplete":true,"isRunning":false, + "isWaiting":false,"reasonForIncompletion":"tool worker never polled"} + """); + + var terminal = Assert.Single(result.Events!); + Assert.Equal(EventType.Error, terminal.Type); + Assert.Equal("tool worker never polled", terminal.Content); + } + + [Fact] + public async Task ToolTaskWithNoOutput_StillReported_WithoutAResult() + { + // A tool the LLM dispatched that failed or never finished records no output. + // Dropping it would hide a call that demonstrably happened. + var result = await WaitWithTasksAsync(""" + {"tasks":[ + {"taskType":"HTTP","taskDefName":"lookup_order","referenceTaskName":"lookup_order_0", + "inputData":{"_agent_tool_name":"lookup_order","order":"A-1"}, + "outputData":{}} + ]} + """); + + var toolCall = Assert.Single(result.ToolCalls!); + Assert.Equal("lookup_order", toolCall["name"]); + Assert.False(toolCall.ContainsKey("result")); + Assert.Equal( + [EventType.ToolCall, EventType.ToolResult, EventType.Done], + result.Events!.Select(e => e.Type)); + } + + [Fact] + public async Task HttpToolResult_FallsBackToWholeOutput() + { + // An HTTP tool answers under `response`, so keying only on `result` loses it. + var result = await WaitWithTasksAsync(""" + {"tasks":[ + {"taskType":"HTTP","taskDefName":"lookup_order","referenceTaskName":"lookup_order_0", + "inputData":{"_agent_tool_name":"lookup_order","order":"A-1"}, + "outputData":{"response":{"status":"shipped"},"statusCode":200}} + ]} + """); + + var toolResult = Assert.IsType(Assert.Single(result.ToolCalls!)["result"]); + Assert.Equal( + ["response", "statusCode"], + toolResult.EnumerateObject().Select(p => p.Name)); + } +} diff --git a/Conductor.AI.Tests/StubAgentServer.cs b/Conductor.AI.Tests/StubAgentServer.cs new file mode 100644 index 00000000..5cf5e4ad --- /dev/null +++ b/Conductor.AI.Tests/StubAgentServer.cs @@ -0,0 +1,65 @@ +/* + * Copyright 2024 Conductor Authors. + *

+ * Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + *

+ * http://www.apache.org/licenses/LICENSE-2.0 + *

+ * Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on + * an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the + * specific language governing permissions and limitations under the License. + */ +// A stubbed agent server for tests that drive AgentHandle against canned +// responses, so an execution's terminal state can be stated as a payload. + +using System.Net; +using System.Net.Http; +using System.Text; +using System.Threading.Tasks; +using Conductor.Client; +using RestSharp; + +namespace Conductor.AI.Tests; + +internal static class StubAgentServer +{ + private sealed class StubHandler : HttpMessageHandler + { + private readonly Func _respond; + public StubHandler(Func respond) => _respond = respond; + + protected override Task SendAsync(HttpRequestMessage request, CancellationToken ct) + { + var (status, body) = _respond(request); + return Task.FromResult(new HttpResponseMessage(status) + { + Content = new StringContent(body, Encoding.UTF8, "application/json"), + }); + } + } + + ///

A client configuration whose every request is answered by . + internal static Configuration Configure(Func respond) + { + var configuration = new Configuration { BasePath = "http://server/api" }; + configuration.ApiClient.RestClient = new RestClient(new RestClientOptions("http://server/api") + { + ConfigureMessageHandler = _ => new StubHandler(respond), + }); + return configuration; + } + + /// + /// Answer the three reads makes on reaching + /// a terminal state: the status, the execution record, and the workflow with its tasks. + /// + internal static (HttpStatusCode, string) Route( + HttpRequestMessage request, string statusBody, string executionBody, string workflowBody) + { + var path = request.RequestUri!.AbsolutePath; + if (path.EndsWith("/status")) return (HttpStatusCode.OK, statusBody); + if (path.Contains("/execution/")) return (HttpStatusCode.OK, executionBody); + return (HttpStatusCode.OK, workflowBody); + } +} diff --git a/Conductor.AI/Result.cs b/Conductor.AI/Result.cs index b69afa92..0cd2c09e 100644 --- a/Conductor.AI/Result.cs +++ b/Conductor.AI/Result.cs @@ -305,8 +305,7 @@ public async Task WaitAsync(CancellationToken cancellationToken = d { // Fetch full execution record for token usage and finish reason var execution = await _http.GetExecutionAsync(_executionId, cancellationToken); - // Java parity: walk the workflow's tasks once to aggregate tool calls - // from call_* tasks (enrichment read — null is fine, yields no tool calls). + // Tool calls and events come from the task list (enrichment read; null is fine). var workflowWithTasks = await _http.GetWorkflowWithTasksAsync(_executionId, cancellationToken); return BuildResult(status!, s, execution, workflowWithTasks); } @@ -536,55 +535,174 @@ private static AgentResult BuildResult( _ => null, }; + var executionId = status["executionId"]?.GetValue() ?? ""; + var error = parsedStatus != Status.Completed + ? status["reasonForIncompletion"]?.GetValue() + : null; + + var (toolCalls, events) = ExtractToolActivity(workflowWithTasks, executionId); + events.Add(parsedStatus == Status.Completed + ? new AgentEvent { Type = EventType.Done, ExecutionId = executionId, Output = outputDict } + : new AgentEvent { Type = EventType.Error, ExecutionId = executionId, Content = error }); + return new AgentResult { - ExecutionId = status["executionId"]?.GetValue() ?? "", + ExecutionId = executionId, Status = parsedStatus, Output = outputDict, - Error = parsedStatus != Status.Completed ? status["reasonForIncompletion"]?.GetValue() : null, - ToolCalls = ExtractToolCalls(workflowWithTasks), + Error = error, + ToolCalls = toolCalls, TokenUsage = tokenUsage, FinishReason = finishReason, + Events = events, }; } + // ── Tool-call extraction ───────────────────────────────────────── + + /// + /// Task types only a tool compiles to: the server's ToolCompiler.TYPE_MAP. + /// GENERATE_PDF is absent from that map, but media tools compile to it. + /// + private static readonly HashSet ToolTaskTypes = new(StringComparer.Ordinal) + { + "HTTP", "CALL_MCP_TOOL", + "GENERATE_IMAGE", "GENERATE_AUDIO", "GENERATE_VIDEO", "GENERATE_PDF", + "LLM_INDEX_TEXT", "LLM_SEARCH_INDEX", "PULL_WORKFLOW_MESSAGES", + }; + + /// Tool task types the agent layer also emits for sub-agents, routers and plan approval. + private static readonly HashSet MarkerRequiredToolTaskTypes = new(StringComparer.Ordinal) + { + "SUB_WORKFLOW", "HUMAN", + }; + + /// Markers the server's dispatch script injects on a tool task it dispatched. + private const string AgentToolNameKey = "_agent_tool_name"; + private const string AgentStateKey = "_agent_state"; + + /// The tool name on the dynamic-tools path, which injects no _agent_tool_name. + private const string MethodKey = "method"; + + /// Runtime keys carried on a tool task's input that are not the tool's arguments. + private static readonly HashSet InternalInputKeys = new(StringComparer.Ordinal) + { + MethodKey, "evaluatorType", "expression", "ctx", "workerTag", "agentConfig", + }; + + /// + /// Whether a task is a tool call the LLM dispatched. Never judged by reference + /// name, which carries the provider's tool-call id. + /// + private static bool IsToolTask(JsonNode task) + { + var taskType = task["taskType"]?.GetValue(); + if (taskType is null) return false; + if (ToolTaskTypes.Contains(taskType)) return true; + + // A worker tool's task type is the tool's own name, so only the marker can + // identify it. The marker also keeps the agent's own structure out: on type + // alone every handoff would read as a tool call that never happened. Cost is + // that the dynamic-tools path marks no non-worker tool, so an agent-as-tool + // dispatched there is missed. + var needsMarker = MarkerRequiredToolTaskTypes.Contains(taskType) + || taskType == task["taskDefName"]?.GetValue(); + return needsMarker + && task["inputData"] is JsonObject inputData + && (inputData.ContainsKey(AgentToolNameKey) || inputData.ContainsKey(AgentStateKey)); + } + + /// + /// The tool's name from the payload, never the task type, which holds the tool + /// name only for a worker tool. + /// + private static string ResolveToolName(JsonNode task) + { + if (task["inputData"] is JsonObject inputData) + { + if (StringValue(inputData, AgentToolNameKey) is { } toolName) return toolName; + if (StringValue(inputData, MethodKey) is { } method) return method; + } + return task["taskDefName"]?.GetValue() ?? ""; + } + + private static string? StringValue(JsonObject obj, string key) + => obj.TryGetPropertyValue(key, out var node) && node is JsonValue value + && value.TryGetValue(out var s) && !string.IsNullOrEmpty(s) + ? s + : null; + + /// The tool's arguments: the task's input with the runtime's own keys stripped. + private static Dictionary ToolArgs(JsonObject inputData) + { + var cleaned = new Dictionary(); + foreach (var kv in inputData) + { + var k = kv.Key; + if (k.StartsWith('_') || InternalInputKeys.Contains(k)) continue; + cleaned[k] = JsonSerializer.Deserialize(kv.Value?.ToJsonString() ?? "null", ConductorAgentJson.Options)!; + } + return cleaned; + } + /// - /// Java parity (AgentHandle.extractFromTasks): walk the workflow's tasks - /// and collect one entry per LLM-dispatched tool call — tasks whose - /// referenceTaskName starts with call_ — capturing the tool name, - /// its input args (internal runtime fields stripped), and its result. + /// The task's result, or its whole output when there is none: an HTTP tool + /// answers under response. Matches the server's AgentEventListener. /// - private static List>? ExtractToolCalls(JsonNode? workflowWithTasks) + private static object? ToolResult(JsonObject outputData) { - if (workflowWithTasks?["tasks"] is not JsonArray tasks) return null; + var node = outputData.TryGetPropertyValue("result", out var resultNode) && resultNode is not null + ? resultNode + : outputData; + return JsonSerializer.Deserialize(node.ToJsonString(), ConductorAgentJson.Options); + } + /// + /// Both views of each tool call in one pass. Reconstructed from finished tasks, + /// so narrower than the live stream. + /// + private static (List>? ToolCalls, List Events) ExtractToolActivity( + JsonNode? workflowWithTasks, string executionId) + { List>? toolCalls = null; + var events = new List(); + + if (workflowWithTasks?["tasks"] is not JsonArray tasks) return (toolCalls, events); + foreach (var task in tasks) { - var refName = task?["referenceTaskName"]?.GetValue(); - var outputData = task?["outputData"] as JsonObject; - if (refName is null || !refName.StartsWith("call_", StringComparison.Ordinal) || outputData is null) - continue; + if (task is null || !IsToolTask(task)) continue; + + var name = ResolveToolName(task); + var args = task["inputData"] is JsonObject inputData ? ToolArgs(inputData) : null; + // A task that recorded no output still ran: report the call without a result + // rather than dropping it, which would hide a tool call that failed. + var result = task["outputData"] is JsonObject { Count: > 0 } outputData + ? ToolResult(outputData) + : null; + + var tc = new Dictionary { ["name"] = name }; + if (args is not null) tc["args"] = args; + if (result is not null) tc["result"] = result; + (toolCalls ??= new List>()).Add(tc); - var tc = new Dictionary { ["name"] = task!["taskType"]?.GetValue() ?? "" }; - if (task["inputData"] is JsonObject inputData) + events.Add(new AgentEvent { - var cleaned = new Dictionary(); - foreach (var kv in inputData) - { - var k = kv.Key; - if (k.StartsWith('_') || k is "method" or "evaluatorType" or "expression" or "ctx" - or "workerTag" or "agentConfig") - continue; - cleaned[k] = JsonSerializer.Deserialize(kv.Value?.ToJsonString() ?? "null", ConductorAgentJson.Options)!; - } - tc["args"] = cleaned; - } - if (outputData.TryGetPropertyValue("result", out var resultNode) && resultNode is not null) - tc["result"] = JsonSerializer.Deserialize(resultNode.ToJsonString(), ConductorAgentJson.Options)!; - - (toolCalls ??= new List>()).Add(tc); + Type = EventType.ToolCall, + ExecutionId = executionId, + ToolName = name, + Args = args, + Timestamp = task["startTime"]?.GetValue(), + }); + events.Add(new AgentEvent + { + Type = EventType.ToolResult, + ExecutionId = executionId, + ToolName = name, + Result = result, + Timestamp = task["endTime"]?.GetValue(), + }); } - return toolCalls; + return (toolCalls, events); } } diff --git a/docs/agents/concepts/streaming-hitl.md b/docs/agents/concepts/streaming-hitl.md index bf0f1b5a..cf653b89 100644 --- a/docs/agents/concepts/streaming-hitl.md +++ b/docs/agents/concepts/streaming-hitl.md @@ -29,6 +29,20 @@ await foreach (var ev in runtime.StreamAsync(agent, "Write a haiku about C#.")) Event types: `Thinking`, `ToolCall`, `ToolResult`, `GuardrailPass`, `GuardrailFail`, `Waiting`, `Handoff`, `Message`, `Error`, `Done`. +### Events on a waited result + +```csharp +var result = await handle.WaitAsync(); +foreach (var ev in result.Events!) Console.WriteLine($"{ev.Type} {ev.ToolName}"); +``` + +`AgentResult.Events` carries the run's tool activity — a `ToolCall`/`ToolResult` pair +per tool call, closed by a terminal `Done`, or `Error` for a run that did not +complete. It is reconstructed from the finished execution's tasks, so it is never +null but is narrower than the live stream: the server also emits `Thinking`, +`Handoff`, guardrail and per-failed-task events, and none of those survives into the +terminal record. Stream the run when the events themselves are the point. + Streaming attempts SSE first and falls back to status-polling. Disable SSE entirely with `CONDUCTOR_AGENT_STREAMING_ENABLED=false` — see [deploy-serve-run.md](deploy-serve-run.md#worker-tuning-and-agentconfig). If the diff --git a/docs/agents/reference/api.md b/docs/agents/reference/api.md index f68439ce..ecac78cf 100644 --- a/docs/agents/reference/api.md +++ b/docs/agents/reference/api.md @@ -211,7 +211,8 @@ Positions map to server task names: `before_agent`, `after_agent`, `before_model **`AgentResult`** (record): `ExecutionId`, `CorrelationId`, `Output` (`Dictionary?`; final text usually `Output["result"]`), `Messages`, -`ToolCalls`, `Status`, `FinishReason`, `Error`, `TokenUsage`, `Metadata`, `Events`, +`ToolCalls`, `Status`, `FinishReason`, `Error`, `TokenUsage`, `Metadata`, +`Events` (see [streaming-hitl.md](../concepts/streaming-hitl.md#events-on-a-waited-result)), `SubResults`. Convenience: `IsSuccess`, `IsFailed`, `IsRejected`, `PrintResult()`. **`AgentHandle`** — `ExecutionId`, `RunId`, `WaitAsync(ct)`, `StreamAsync(ct)`,