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 @@ -572,7 +572,7 @@ internal static string MissingPayloadRepairPrompt(string kind, string defect, st
/// <summary>The fewest decisions worth folding — below this a compaction would not shrink the prompt meaningfully (the overflow has another cause), so the original fault propagates to the clean-stop path.</summary>
internal const int MinCompactFold = 4;

private static readonly JsonElement TapeSummarySchema = JsonDocument.Parse("""
internal static readonly JsonElement TapeSummarySchema = JsonDocument.Parse("""
{ "type": "object", "additionalProperties": false, "required": ["summary"], "properties": { "summary": { "type": "string", "description": "The compact progress digest." } } }
""").RootElement;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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.JsonSchema } },
Tools = new[] { new AnthropicTool { Name = StructuredToolName, Description = "Return the result as structured JSON.", InputSchema = request.WireJsonSchema ?? request.JsonSchema } },
ToolChoice = new AnthropicToolChoice { Type = "tool", Name = StructuredToolName }
};

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,15 @@ public sealed record StructuredLLMCompletionRequest
/// <summary>JSON Schema (object) the response MUST conform to.</summary>
public required JsonElement JsonSchema { get; init; }

/// <summary>
/// The schema the PROVIDER is handed as the forced tool's schema, when that must differ from <see cref="JsonSchema"/>;
/// null sends <see cref="JsonSchema"/> itself. It exists for a constrained decoder that cannot compile a construct the
/// contract needs (see <see cref="JsonSchemaCombinators"/>). It must ACCEPT everything <see cref="JsonSchema"/>
/// accepts: every reply is still validated against <see cref="JsonSchema"/>, which also stays the schema the prompt
/// quotes, so a narrower wire schema would forbid the model an answer the contract allows.
/// </summary>
public JsonElement? WireJsonSchema { get; init; }

/// <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
@@ -0,0 +1,63 @@
using System.Text.Json;
using System.Text.Json.Nodes;

namespace CodeSpace.Core.Services.Workflows.Llm;

/// <summary>
/// Derives a schema's COMBINATOR-FREE form: <c>oneOf</c>, <c>anyOf</c>, <c>allOf</c>, <c>not</c>, <c>if</c>, <c>then</c>
/// and <c>else</c> dropped at every schema position, every other keyword, property name and literal kept verbatim. This
/// is the form a provider's constrained decoder is handed as <see cref="StructuredLLMCompletionRequest.WireJsonSchema"/>
/// while the full schema keeps validating every reply.
///
/// <para>Why it exists: a hosted vLLM backend compiles the forced tool's schema into a decoding grammar, and a
/// combinator shape it cannot compile fails the whole call with an EMPTY HTTP 500 — not a 400 the progressive fallback
/// degrades on. The planner's per-kind acceptance branches did exactly that on every call, while every combinator-free
/// schema sent to the same model was answered. The error names no keyword, so every combinator goes, not a guessed one.</para>
///
/// <para>Why it is safe: each combinator only adds a constraint alongside its siblings, so dropping one can only widen
/// what a schema accepts, and the model is never forbidden a reply the contract allows. What the wire no longer says is
/// still enforced: the bounded re-ask validates against the full schema. The one keyword family this reasoning would
/// not hold for, <c>unevaluatedProperties</c>/<c>unevaluatedItems</c>, reads annotations out of these subschemas; no
/// model-facing schema here uses it.</para>
/// </summary>
public static class JsonSchemaCombinators
{
private static readonly string[] Keywords = ["oneOf", "anyOf", "allOf", "not", "if", "then", "else"];

/// <summary>Keywords whose value maps NAMES to subschemas — the names are data (a property called <c>not</c> is still a property), only the values are schemas.</summary>
private static readonly HashSet<string> NamedSubschemas = ["properties", "patternProperties", "$defs", "definitions", "dependentSchemas"];

/// <summary>Keywords whose value is a subschema, or an array of them. Anything else (<c>enum</c>, <c>const</c>, <c>default</c>, <c>required</c>, …) is data and is copied untouched.</summary>
private static readonly HashSet<string> Subschemas = ["items", "prefixItems", "additionalItems", "additionalProperties", "contains", "propertyNames", "unevaluatedItems", "unevaluatedProperties"];

/// <summary>A copy of <paramref name="schema"/> with every combinator removed; the input is never modified.</summary>
public static JsonElement Strip(JsonElement schema)
{
var copy = JsonNode.Parse(schema.GetRawText());

StripSchema(copy);

return JsonSerializer.SerializeToElement(copy);
}

private static void StripSchema(JsonNode? node)
{
if (node is not JsonObject schema) return;

foreach (var keyword in Keywords) schema.Remove(keyword);

foreach (var (keyword, value) in schema) StripChildren(keyword, value);
}

private static void StripChildren(string keyword, JsonNode? value)
{
if (NamedSubschemas.Contains(keyword) && value is JsonObject named) StripEach(named.Select(pair => pair.Value));
else if (Subschemas.Contains(keyword) && value is JsonArray list) StripEach(list);
else if (Subschemas.Contains(keyword)) StripSchema(value);
}

private static void StripEach(IEnumerable<JsonNode?> schemas)
{
foreach (var schema in schemas) StripSchema(schema);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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.JsonSchema } } },
Tools = new[] { new OpenAiTool { Function = new OpenAiFunction { Name = StructuredToolName, Description = "Return the result as structured JSON.", Parameters = request.WireJsonSchema ?? request.JsonSchema } } },
ToolChoice = new OpenAiToolChoice { Function = new OpenAiToolChoiceFunction { Name = StructuredToolName } },
};

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -54,7 +54,7 @@ public static NodeResult Park(NodeRunContext context, LlmApiException fault, Dat
var delay = SupervisorInfraPark.DelayFor(state.Parks);
var marker = SupervisorInfraPark.Marker(state, 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", context.NodeId, fault.Category, delay, state.Parks, state.FirstParkedAtUtc);
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);

return NodeResult.Suspend(new SuspensionToken
{
Expand Down
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
using System.Text.Json;
using System.Text.Json.Serialization;
using CodeSpace.Core.Services.Workflows.Llm;

namespace CodeSpace.Core.Services.Workflows.Planning;

Expand Down Expand Up @@ -155,6 +156,15 @@ public static class PlannerSchema
}
""").RootElement.Clone();

/// <summary>
/// What the provider's constrained decoder is handed in place of <see cref="ResponseSchema"/>: the same schema with its
/// combinators stripped, so the acceptance is one flat typed object that still declares every field (<c>formatVersion</c>,
/// <c>kind</c>, <c>argv</c>, <c>artifactPaths</c>, …) but carries no per-kind branch. The branches made a hosted vLLM
/// backend answer every planner call with an empty HTTP 500. <see cref="ResponseSchema"/> still validates each reply and
/// is still the schema the prompt quotes, so the per-kind requirements are enforced by the bounded re-ask instead.
/// </summary>
public static readonly JsonElement WireSchema = JsonSchemaCombinators.Strip(ResponseSchema);

/// <summary>Deserialization options for mapping a schema-valid object back into <c>PlannedWorkflow</c>. Case-insensitive so the model's lower-camel keys bind to the record's Pascal properties; the string-enum converter binds the acceptance <c>kind</c> ("TestsPass"/"ArtifactPresent") to <c>BenchmarkGradingKind</c>.</summary>
public static readonly JsonSerializerOptions Options = new()
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -91,6 +91,7 @@ public async Task<PlannedWorkflow> PlanAsync(WorkflowPlanRequest request, Cancel
SystemPrompt = SystemPrompt,
UserPrompt = BuildUserPrompt(request, catalog, lessons),
JsonSchema = PlannerSchema.ResponseSchema,
WireJsonSchema = PlannerSchema.WireSchema,
ResponseValidator = ValidateModelResponse,
ResponseAdvisor = AdviseModelResponse,
MaxOutputTokens = 4096,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -62,6 +62,7 @@ public static async Task<StructuredLLMCompletionRequest> BuildAsync(ILifetimeSco
SystemPrompt = LlmWorkflowPlanner.SystemPrompt,
UserPrompt = LlmWorkflowPlanner.BuildUserPromptForTest(planRequest, catalog),
JsonSchema = PlannerSchema.ResponseSchema,
WireJsonSchema = PlannerSchema.WireSchema,
};
}

Expand Down
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
using System.Diagnostics;
using System.Text.Json;
using Autofac;
using CodeSpace.Core.Persistence.Db;
using CodeSpace.Core.Services.Workflows.Engine;
Expand Down Expand Up @@ -159,7 +160,18 @@ private static async Task<TimeSpan> PauseAsync(Stopwatch clock, TimeSpan wakePau

private static string Unresolved(ParkedCell cell, int wakes, TimeSpan wakePause) =>
$"node '{cell.NodeId}' was STILL parked on {WorkflowWaitKinds.SupervisorInfraPark} after {wakes} deadline wake(s) over ~{(wakes * wakePause).TotalSeconds:0}s "
+ "— the model plane never came back inside the ride's budget. That is INFRA, not a model verdict, and it is NOT a pass: nothing was driven to completion.";
+ "— the model plane never came back inside the ride's budget. That is INFRA, not a model verdict, and it is NOT a pass: nothing was driven to completion. "
+ $"The park's last fault: {ParkFault(cell)}";

/// <summary>The fault the park stored on its own marker — the gateway's own words — so a park that outlives the ride names its cause in the job summary, not only its duration.</summary>
private static string ParkFault(ParkedCell cell)
{
if (cell.WaitPayloadJson is null) return "(none recorded)";

using var marker = JsonDocument.Parse(cell.WaitPayloadJson);

return marker.RootElement.TryGetProperty("error", out var error) && error.ValueKind == JsonValueKind.String ? error.GetString()! : "(none recorded)";
}

/// <summary>Read the run's newest pending infra park plus the projected status of the cell it holds. No park pending ⇒ a Settled cell (there is nothing for the ride to wake).</summary>
private static async Task<ParkedCell> ReadCellAsync(PostgresFixture fixture, Guid runId)
Expand Down
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
using CodeSpace.Core.Services.Supervisor;
using CodeSpace.Messages.Constants;
using CodeSpace.Messages.Enums;
using Shouldly;
Expand Down Expand Up @@ -122,6 +123,20 @@ public async Task An_exhausted_budget_yields_the_honest_infra_outcome_never_a_gr
RealModelGate.IsGatewayInfraFailure(new AggregateException(ex)).ShouldBeTrue("the await chain can wrap it");
}

[Fact]
public async Task An_unresolved_park_names_the_fault_the_park_stored_on_its_marker()
{
// The planner's park outlived this ride on run after run, and the skip said only THAT the model plane never
// came back. What the gateway actually said sat on the park's own marker, one read away. The marker is minted
// by the production writer, so the key this read depends on cannot drift from it unnoticed.
var marker = SupervisorInfraPark.Marker(SupervisorInfraPark.Next(null, DateTimeOffset.UtcNow), "Anthropic API error (HTTP 500, Transient): Hosted_vllmException");
var cell = Parked() with { WaitPayloadJson = marker.GetRawText() };

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

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 void The_ride_pauses_for_real_between_wakes_so_a_recovery_can_actually_be_observed()
{
Expand Down
25 changes: 25 additions & 0 deletions backend/tests/CodeSpace.UnitTests/Workflows/InfraParkTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
using CodeSpace.Core.Services.Workflows.Runtime;
using CodeSpace.Messages.Constants;
using CodeSpace.Messages.Enums;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Logging.Abstractions;
using Shouldly;

Expand Down Expand Up @@ -63,6 +64,18 @@ public void A_first_fault_parks_on_the_shared_wait_kind_with_a_deadline_wake()
result.SuspendUntil.TimeoutPayload.ShouldNotBeNull("the wake must carry the ladder position forward");
}

[Fact]
public void The_park_log_carries_the_faults_own_words()
{
// The planner parked on an empty-bodied gateway 500 run after run, and the park line named only the category —
// so the one thing a reader needed, what the gateway actually said, was in no log line at all.
var logger = new CapturingLogger();

InfraPark.Park(Context() with { Logger = logger }, Fault(LlmErrorCategory.Transient), DateTimeOffset.UtcNow);

logger.Messages.ShouldHaveSingleItem().ShouldContain("upstream unavailable");
}

[Fact]
public void The_park_keeps_the_nodes_ambient_cell_so_a_map_branch_stays_in_its_branch()
{
Expand Down Expand Up @@ -113,4 +126,16 @@ public void Past_the_whole_window_the_node_fails_honestly_instead_of_parking_for
result.Error.ShouldContain("model plane", Case.Insensitive);
result.Retryable.ShouldBeFalse("re-running the node cannot reach a provider that has been down for a day");
}

/// <summary>Keeps every line the park writes, formatted the way a sink would render it.</summary>
private sealed class CapturingLogger : ILogger
{
public List<string> Messages { get; } = [];

public IDisposable BeginScope<TState>(TState state) where TState : notnull => NullScope.Instance;
public bool IsEnabled(LogLevel logLevel) => true;
public void Log<TState>(LogLevel logLevel, EventId eventId, TState state, Exception? exception, Func<TState, Exception?, string> formatter) => Messages.Add(formatter(state, exception));

private sealed class NullScope : IDisposable { public static readonly NullScope Instance = new(); public void Dispose() { } }
}
}
Loading
Loading