Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -14,8 +14,8 @@ namespace CodeSpace.Core.Services.Workflows.Llm.Anthropic;
///
/// Implements <see cref="ILLMClient"/> (free-text completion) AND the sibling
/// <see cref="IStructuredLLMClient"/> (schema-constrained JSON) — the latter via Anthropic's
/// forced tool-use: a single tool whose <c>input_schema</c> IS the requested schema, with
/// <c>tool_choice</c> pinned to it, so the model's <c>tool_use</c> block carries schema-valid JSON.
/// forced tool-use: a single tool whose <c>input_schema</c> is the request's <c>WireJsonSchema ?? JsonSchema</c> (every reply
/// is validated against <c>JsonSchema</c>), with <c>tool_choice</c> pinned to it, so the model's <c>tool_use</c> block carries schema-valid JSON.
///
/// Keep ONLY the wire-shape concerns here. Anything node-facing (prompt assembly, output
/// trimming, retry policy) belongs in the llm.complete node or the LLM-side resilience
Expand Down Expand Up @@ -195,7 +195,7 @@ private async Task<StructuredLLMCompletion> CompleteStructuredOnceAsync(Structur
}

/// <summary>
/// Attempt 1: forced tool-use — a single tool whose input_schema IS the schema, tool_choice pinned to it. Returns
/// Attempt 1: forced tool-use — a single tool whose input_schema is <c>WireJsonSchema ?? JsonSchema</c> (validation still uses <c>JsonSchema</c>), tool_choice pinned to it. Returns
/// the recovered JSON (or null to degrade to the prompt-only floor) PLUS the parsed response — so the caller can
/// accumulate the BILLED usage even when the JSON is null (a 200 that produced no usable tool call still cost
/// tokens). On a 400/422 reject the request was never generated, so both are null (nothing to bill).
Expand All @@ -213,7 +213,7 @@ private async Task<StructuredLLMCompletion> CompleteStructuredOnceAsync(Structur
StopSequences = request.Sampling?.Stop,
System = system,
Messages = messages,
Tools = new[] { new AnthropicTool { Name = StructuredToolName, Description = "Return the result as structured JSON.", InputSchema = request.WireJsonSchema ?? request.JsonSchema } },
Tools = new[] { new AnthropicTool { Name = StructuredToolName, Description = "Return the result as structured JSON.", InputSchema = request.ProviderSchema } },
ToolChoice = new AnthropicToolChoice { Type = "tool", Name = StructuredToolName }
};

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,9 @@ public interface IStructuredLLMClient
/// <summary>
/// One LLM call constrained to return JSON matching <see cref="StructuredLLMCompletionRequest.JsonSchema"/>.
/// The provider impl maps the schema to whatever its API offers (Anthropic forces a single tool whose
/// input_schema IS the schema; OpenAI would use response_format json_schema, …).
/// input_schema carries the schema; OpenAI would use response_format json_schema, …). The provider receives
/// <see cref="StructuredLLMCompletionRequest.WireJsonSchema"/> when the request sets one, else <see cref="StructuredLLMCompletionRequest.JsonSchema"/>;
/// every reply is validated against <see cref="StructuredLLMCompletionRequest.JsonSchema"/> either way.
/// </summary>
Task<StructuredLLMCompletion> CompleteStructuredAsync(StructuredLLMCompletionRequest request, CancellationToken cancellationToken);
}
Expand All @@ -45,6 +47,9 @@ public sealed record StructuredLLMCompletionRequest
/// </summary>
public JsonElement? WireJsonSchema { get; init; }

/// <summary>The schema the provider is actually handed: <see cref="WireJsonSchema"/> when set, else <see cref="JsonSchema"/>. A <c>default</c> wire schema (kind Undefined) counts as unset — a nullable struct admits it as a non-null value, and it cannot be serialized.</summary>
internal JsonElement ProviderSchema => WireJsonSchema is { ValueKind: not JsonValueKind.Undefined } wire ? wire : JsonSchema;

/// <summary>Server-only validation of the consumer contract, alongside JSON schema. Violations enter the same bounded model re-ask; this callback never rewrites output.</summary>
[JsonIgnore]
public Func<JsonElement, IReadOnlyList<string>>? ResponseValidator { get; init; }
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,8 +18,8 @@ namespace CodeSpace.Core.Services.Workflows.Llm.OpenAi;
/// post-S6b) — a call without a credential fails closed.
///
/// Implements <see cref="ILLMClient"/> (free-text) AND <see cref="IStructuredLLMClient"/> (schema-constrained
/// JSON). Structured output uses FORCED FUNCTION-CALLING — a single function whose <c>parameters</c> IS the
/// requested schema, with <c>tool_choice</c> pinned to it — rather than <c>response_format: json_schema</c>,
/// JSON). Structured output uses FORCED FUNCTION-CALLING — a single function whose <c>parameters</c> is the request's
/// <c>WireJsonSchema ?? JsonSchema</c> (every reply is validated against <c>JsonSchema</c>), with <c>tool_choice</c> pinned to it — rather than <c>response_format: json_schema</c>,
/// because function-calling is supported by far more OpenAI-compatible gateways than the newer structured-outputs
/// feature, and it mirrors the Anthropic client's forced-tool design exactly (one coercion mechanism to reason
/// about). The model's <c>tool_calls[0].function.arguments</c> is the schema-SHAPED JSON (classic function-calling
Expand Down Expand Up @@ -229,7 +229,7 @@ private async Task<StructuredLLMCompletion> CompleteStructuredOnceAsync(Structur
}

/// <summary>
/// Attempt 1: forced function-calling — a single function whose parameters IS the schema, tool_choice pinned to it.
/// Attempt 1: forced function-calling — a single function whose parameters is <c>WireJsonSchema ?? JsonSchema</c> (validation still uses <c>JsonSchema</c>), tool_choice pinned to it.
/// Returns the recovered JSON (or null to degrade to the prompt-only floor) PLUS the parsed response — so the caller
/// can accumulate the BILLED usage even when the JSON is null (a 200 that produced no usable function call still
/// cost tokens). On a 400/422 reject the request was never generated, so both are null (nothing to bill).
Expand All @@ -254,7 +254,7 @@ private async Task<StructuredLLMCompletion> CompleteStructuredOnceAsync(Structur
Stop = request.Sampling?.Stop,
ReasoningEffort = LlmModelCapabilities.SupportsReasoningEffort(request.Model) ? request.ReasoningEffort : null, // sent ONLY to a reasoning model (a plain chat model 400s on it); the value rides verbatim (the API validates it per model)
Messages = BuildMessages(system, request.UserPrompt),
Tools = new[] { new OpenAiTool { Function = new OpenAiFunction { Name = StructuredToolName, Description = "Return the result as structured JSON.", Parameters = request.WireJsonSchema ?? request.JsonSchema } } },
Tools = new[] { new OpenAiTool { Function = new OpenAiFunction { Name = StructuredToolName, Description = "Return the result as structured JSON.", Parameters = request.ProviderSchema } } },
ToolChoice = new OpenAiToolChoice { Function = new OpenAiToolChoiceFunction { Name = StructuredToolName } },
};

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -230,9 +230,10 @@ private async Task<NodeResult> ParkForInfraOrStopAsync(NodeRunContext context, G
}

var delay = SupervisorInfraPark.DelayFor(state.Parks);
var marker = SupervisorInfraPark.Marker(state, fault.Message);
var said = InfraPark.FaultText(context.Scope, fault);
var marker = SupervisorInfraPark.Marker(state, said);

context.Logger.LogWarning("agent.supervisor run {RunId}: brain call hit a {Category} infra fault — parking {Delay} (park {Parks} since {First:o}) instead of failing the run", supervisorRunId, fault.Category, delay, state.Parks, state.FirstParkedAtUtc);
context.Logger.LogWarning("agent.supervisor run {RunId}: brain call hit a {Category} infra fault — parking {Delay} (park {Parks} since {First:o}) instead of failing the run: {Fault}", supervisorRunId, fault.Category, delay, state.Parks, state.FirstParkedAtUtc, said);

return NodeResult.Suspend(new SuspensionToken
{
Expand Down
35 changes: 32 additions & 3 deletions backend/src/CodeSpace.Core/Services/Workflows/Nodes/InfraPark.cs
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
using CodeSpace.Core.Services.Supervisor;
using CodeSpace.Core.Services.Workflows.Llm;
using CodeSpace.Core.Services.Workflows.Runtime;
using CodeSpace.Messages.Constants;
using Microsoft.Extensions.Logging;

Expand Down Expand Up @@ -43,18 +44,19 @@ public static class InfraPark
public static NodeResult Park(NodeRunContext context, LlmApiException fault, DateTimeOffset now)
{
var state = SupervisorInfraPark.Next(context.ResumePayload, now);
var said = FaultText(context.Scope, fault);

if (state.WindowExhausted)
{
context.Logger.LogWarning("Node {NodeId}: the model plane stayed unavailable past the whole {Window} park window — failing the node honestly", context.NodeId, SupervisorInfraPark.MaxParkWindow);

return NodeResult.Fail($"The model plane stayed unavailable for {SupervisorInfraPark.MaxParkWindow.TotalHours:0}h: {fault.Message}", retryable: false);
return NodeResult.Fail($"The model plane stayed unavailable for {SupervisorInfraPark.MaxParkWindow.TotalHours:0}h: {said}", retryable: false);
}

var delay = SupervisorInfraPark.DelayFor(state.Parks);
var marker = SupervisorInfraPark.Marker(state, fault.Message);
var marker = SupervisorInfraPark.Marker(state, said);

context.Logger.LogWarning("Node {NodeId}: model call hit a {Category} infra fault — parking {Delay} (park {Parks} since {First:o}) instead of failing the run: {Fault}", context.NodeId, fault.Category, delay, state.Parks, state.FirstParkedAtUtc, fault.Message);
context.Logger.LogWarning("Node {NodeId}: model call hit a {Category} infra fault — parking {Delay} (park {Parks} since {First:o}) instead of failing the run: {Fault}", context.NodeId, fault.Category, delay, state.Parks, state.FirstParkedAtUtc, said);

return NodeResult.Suspend(new SuspensionToken
{
Expand All @@ -68,4 +70,31 @@ public static NodeResult Park(NodeRunContext context, LlmApiException fault, Dat
TimeoutPayload = marker,
});
}

/// <summary>The most characters of a fault's text a park writes down. The gateway's own words are what a reader needs; the provider's whole error body, which <see cref="Exception.Message"/> ends with, is not.</summary>
internal const int MaxFaultChars = 512;

/// <summary>
/// The fault's own words in the form everything that OUTLIVES the call may hold — the park marker (durable, and what the
/// run detail shows while parked), the honest failure text and the park log all take this, never the raw
/// <see cref="Exception.Message"/>, which ends with the provider's error body verbatim and unbounded. Scope secrets are
/// redacted FIRST and the result is clamped to <see cref="MaxFaultChars"/> after: a clamp that ran first could cut a
/// secret in half and leave a fragment no redactor would recognise.
/// </summary>
internal static string FaultText(NodeRunScope scope, LlmApiException fault)
{
var redacted = PersistenceSecretRedactor.FromScope(scope).Redact(fault.Message).Value ?? PersistenceSecretRedactor.Marker;

return Clamp(redacted, MaxFaultChars);
}

/// <summary>The first <paramref name="max"/> characters plus an ellipsis, never ending inside a surrogate pair — a lone half is ill-formed text that Npgsql's strict UTF-8 encoder refuses to write.</summary>
private static string Clamp(string text, int max)
{
if (text.Length <= max) return text;

var cut = char.IsHighSurrogate(text[max - 1]) ? max - 1 : max;

return text[..cut] + "…";
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -5,8 +5,8 @@
namespace CodeSpace.Core.Services.Workflows.Planning;

/// <summary>
/// The planner's COMMIT-CONTRACT: the JSON Schema the model is constrained to (via the structured-output
/// path) and the matching deserialization options. Co-located with the planner concern (Rule 18) and
/// The planner's COMMIT-CONTRACT: the JSON Schema every reply is VALIDATED against (and the prompt quotes) and the
/// matching deserialization options; the provider itself is handed the combinator-free <see cref="WireSchema"/> instead. Co-located with the planner concern (Rule 18) and
/// pinned by a unit test — a drift in either the schema or the property mapping is a contract change a
/// reviewer must see, not an invisible refactor.
///
Expand All @@ -18,7 +18,7 @@ namespace CodeSpace.Core.Services.Workflows.Planning;
/// </summary>
public static class PlannerSchema
{
/// <summary>The JSON schema constraining fresh model output. PlannerAcceptanceDraft maps its typed acceptance payloads before the normalized PlannedWorkflow reaches persistence or execution.</summary>
/// <summary>The JSON schema every fresh model reply is VALIDATED against and the prompt quotes; the provider is handed <see cref="WireSchema"/>, not this. PlannerAcceptanceDraft maps its typed acceptance payloads before the normalized PlannedWorkflow reaches persistence or execution.</summary>
public static readonly JsonElement ResponseSchema = JsonDocument.Parse("""
{
"type": "object",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -13,8 +13,8 @@ namespace CodeSpace.Core.Services.Workflows.Planning.Planners;
/// <summary>
/// The structured-LLM <see cref="IWorkflowPlanner"/> (Rule 18.3 — an impl in the <c>Planners/</c> variant
/// folder). It resolves a structured-capable LLM client through the SAME <see cref="ILLMClientRegistry"/>
/// the <c>llm.complete</c> node uses, sends a system+user prompt constrained by
/// <see cref="PlannerSchema.ResponseSchema"/>, and deserializes the schema-valid object into a
/// the <c>llm.complete</c> node uses, sends a system+user prompt whose reply is validated against
/// <see cref="PlannerSchema.ResponseSchema"/> (the provider is handed the combinator-free <see cref="PlannerSchema.WireSchema"/>), and deserializes the schema-valid object into a
/// <see cref="PlannedWorkflow"/>. Fails cleanly when no registered provider offers structured output.
///
/// <para>The planner produces DATA only — it never wires nodes or runs anything. The grounding context
Expand Down
Original file line number Diff line number Diff line change
@@ -1,6 +1,11 @@
using System.Text.Json;
using CodeSpace.Core.Services.Supervisor;
using CodeSpace.Core.Services.Workflows.Llm;
using CodeSpace.Core.Services.Workflows.Nodes;
using CodeSpace.Core.Services.Workflows.Runtime;
using CodeSpace.Messages.Constants;
using CodeSpace.Messages.Enums;
using Microsoft.Extensions.Logging.Abstractions;
using Shouldly;

namespace CodeSpace.IntegrationTests.Workflows.Supervisor;
Expand Down Expand Up @@ -137,6 +142,21 @@ public async Task An_unresolved_park_names_the_fault_the_park_stored_on_its_mark
ex.Message.ShouldContain("Anthropic API error (HTTP 500, Transient): Hosted_vllmException", Case.Sensitive, "the infra skip must carry the gateway's own words, not only the park's duration");
}

[Fact]
public async Task An_unresolved_park_quotes_the_clamped_text_the_production_park_stored()
{
// The park clamps the gateway's words before it stores them, so a provider that answers with a whole page no
// longer floods the job summary: the skip quotes the first 512 characters and stops. The marker is minted by the
// production park (InfraPark.Park) and read back by the ride's own reader, not built by hand.
var fault = new LlmApiException("Anthropic", 500, LlmErrorCategory.Transient, new string('x', 5_000));
var park = InfraPark.Park(ParkingContext(), fault, DateTimeOffset.UtcNow);
var cell = Parked() with { WaitPayloadJson = park.SuspendUntil!.Payload.GetRawText() };

var ex = await Should.ThrowAsync<InfraParkUnresolvedException>(() => InfraParkRide.RideAsync(() => Task.FromResult(cell), _ => Task.CompletedTask, maxWakes: 1, wakePause: TimeSpan.Zero));

ex.Message.ShouldEndWith("The park's last fault: " + fault.Message[..512] + "…", Case.Sensitive);
}

[Fact]
public void The_ride_pauses_for_real_between_wakes_so_a_recovery_can_actually_be_observed()
{
Expand All @@ -149,4 +169,16 @@ public void The_ride_pauses_for_real_between_wakes_so_a_recovery_can_actually_be
private static ParkedCell Parked() => new() { CellStatus = NodeStatus.Suspended, PendingWaitKind = WorkflowWaitKinds.SupervisorInfraPark, NodeId = "planner" };

private static ParkedCell Settled() => new() { CellStatus = NodeStatus.Success, PendingWaitKind = null, NodeId = "planner" };

private static NodeRunContext ParkingContext() => new()
{
Inputs = new Dictionary<string, JsonElement>(),
Config = new Dictionary<string, JsonElement>(),
RawInputs = JsonDocument.Parse("{}").RootElement,
RawConfig = JsonDocument.Parse("{}").RootElement,
Scope = new NodeRunScope { Trigger = new Dictionary<string, JsonElement>(), Sys = new Dictionary<string, JsonElement>() },
Logger = NullLogger.Instance,
Observability = NodeObservability.NoOp,
NodeId = "planner",
};
}
Loading
Loading