Skip to content
Open
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 @@ -253,17 +253,34 @@ private static JsonElement SynthConfig(TaskBuildContext context)
/// is the identical serialization of the identical array, so an ordinary run's prompt does not change by a
/// character; over budget the model is handed a fair share of every included branch and TOLD, in the prompt's
/// first sentence, that it is reading an excerpt.</para>
///
/// <para>A repo-bound graph's reduce is ALSO handed the integrate node's own account of what landed
/// (<see cref="IntegrationOutcomeSection"/>) and told what to do with it (<see cref="SynthIntegrationInstruction"/>).
/// The integrate step's outputs used to reach only <see cref="DoneInputs"/>: nothing the reduce reads said that a
/// candidate conflicted, or that a unit whose own check rejected it was withheld from it, so a partial integration
/// was narrated as a whole deliverable. Whether the graph integrates is a BUILD-time fact (the integrate node exists
/// or it does not), so the conditional is legitimate here — and a repo-less graph keeps both halves of its prompt
/// byte-for-byte. The account itself is a RUN-time fact the node renders; the prompt only binds it. It sits OUTSIDE
/// the map's <c>promptBudgetChars</c>, which bounds the results projection alone — the node bounds the account
/// instead (<c>RunIntegrationSummary.MaxChars</c>, 2,000 characters against a 120,000-character default budget).</para>
/// </summary>
private static JsonElement SynthInputs(TaskBuildContext context) => JsonSerializer.SerializeToElement(new
private static JsonElement SynthInputs(TaskBuildContext context)
{
systemPrompt = SynthSystemPrompt,
// BYTE-IDENTICAL to the pre-continue prompt: the data half of the reduce carries the goal and the results,
// and nothing else. The failure half rides the SYSTEM prompt instead, so a run in which nothing failed is
// never handed a "Subtasks that failed: 0" line — the reduce learns about a failure from the marker actually
// sitting in its results, which is the only place the fact exists per-run. (A build-time conditional cannot
// express this: the failure count is a RUN-time fact and the prompt is frozen into the definition.)
userPrompt = $"Goal: {context.Seed.Goal}\n\nPer-subtask results:\n" + SynthResultsRef,
});
var integrates = context.AgentProfile?.RepositoryId is not null;

return JsonSerializer.SerializeToElement(new
{
systemPrompt = integrates ? SynthSystemPrompt + SynthIntegrationInstruction : SynthSystemPrompt,
// BYTE-IDENTICAL to the pre-continue prompt on a repo-less graph: the data half of the reduce carries the
// goal and the results, and nothing else. The failure half rides the SYSTEM prompt instead, so a run in
// which nothing failed is never handed a "Subtasks that failed: 0" line — the reduce learns about a failure
// from the marker actually sitting in its results, which is the only place the fact exists per-run. (A
// build-time conditional cannot express this: the failure count is a RUN-time fact and the prompt is frozen
// into the definition.) A repo-bound graph also carries what landed — on a clean run ONE factual sentence —
// because that is the fact the reduce needs in order to call the work delivered, not failure furniture.
userPrompt = $"Goal: {context.Seed.Goal}\n\nPer-subtask results:\n" + SynthResultsRef + (integrates ? IntegrationOutcomeSection : ""),
});
}

/// <summary>
/// The reduce's instruction. The failure clause is what continue-on-error requires of it: the run now reaches
Expand All @@ -278,6 +295,19 @@ private static JsonElement SynthInputs(TaskBuildContext context) => JsonSerializ
+ "A subtask that FAILED appears in the results as an {\"error\": ...} entry instead of a result: never present its work as done — "
+ "say which subtasks failed, what they were meant to deliver, and what is therefore missing from the answer.";

/// <summary>
/// What a REPO-BOUND reduce is told about the integration outcome it is shown, appended to
/// <see cref="SynthSystemPrompt"/> only when the graph has an integrate node — a repo-less reduce has no such
/// outcome, and a sentence about one would only invite the model to invent it. The two verbs carry the honesty
/// contract: what landed is stated as delivered, and anything conflicted or withheld is named as NOT delivered.
/// </summary>
internal const string SynthIntegrationInstruction =
" The integration outcome tells you what actually landed on the integrated branch; state what landed, "
+ "and name anything conflicted or withheld as NOT delivered — never narrate withheld work as done.";

/// <summary>The data half's integration section (repo-bound graphs only): the integrate node's own rendering of what landed, bound whole into the prompt. The key is one <c>git.integrate_run</c> declares in its OutputSchema, which <c>DefinitionValidator</c> enforces at build.</summary>
private const string IntegrationOutcomeSection = "\n\nIntegration outcome:\n{{nodes.integrate.outputs.summary}}";

/// <summary>The reduce's results binding, composed from <see cref="WorkflowOutputKeys.MapResultsPrompt"/> so the prompt and the key the reducer writes cannot drift apart.</summary>
private const string SynthResultsRef = "{{nodes.map.outputs." + WorkflowOutputKeys.MapResultsPrompt + "}}";

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -156,12 +156,14 @@ private static string BuildIntegrationBranch(NodeRunContext context) =>
["appliedCount"] = JsonSerializer.SerializeToElement(result.AppliedCount),
["reason"] = JsonSerializer.SerializeToElement(result.Reason),
["conflicts"] = JsonSerializer.SerializeToElement(
result.Outcomes
.Where(o => o.Disposition != ContributionDisposition.Applied)
.OrderBy(o => o.Skipped)
NotApplied(result)
.Select(o => new { label = o.Label, disposition = o.Disposition.ToString(), reason = o.Reason, conflictedFiles = o.ConflictedFiles, fallbackBranch = o.FallbackBranch, skipped = o.Skipped })),
};

/// <summary>The outcomes that did NOT apply, a real failure before a bystander (stable — ties keep outcome order). The one ordering <c>conflicts[]</c> above and <c>git.integrate_run</c>'s prose summary both read, so what a reader is told first never depends on which of the two it reads.</summary>
internal static IEnumerable<ContributionOutcome> NotApplied(IntegrationResult result) =>
result.Outcomes.Where(o => o.Disposition != ContributionDisposition.Applied).OrderBy(o => o.Skipped);

private static IReadOnlyList<BranchContribution> ReadContributions(NodeRunContext context)
{
if (!context.Inputs.TryGetValue("contributions", out var value) || value.ValueKind != JsonValueKind.Array)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,13 @@ namespace CodeSpace.Core.Services.Workflows.Nodes.Builtin;
/// resumed pass re-derives and RE-INTEGRATES against then-current facts (the human may have pushed a fix or a
/// reconciled branch), so an approve after a repair lands the Clean candidate; a still-conflicted retry
/// completes honestly with the review trail on its outputs — one park per run, never a loop.</para>
///
/// <para>The node owns the NARRATIVE of its own outcome. Every pass — clean, conflicted, skipped, resumed — emits
/// <c>summary</c> (<see cref="RunIntegrationSummary"/>: what actually landed, what conflicted and where its work is
/// kept) and <c>withheld</c> (<see cref="RunIntegrationContributions.Withheld"/>: the units the head gate kept off the
/// candidate because their own definition-of-done rejected them — dropped BEFORE integration, so in no outcome).
/// The plan-map synth is handed the summary: before it, the reduce narrated a partial integration as a whole
/// deliverable, because <c>appliedCount</c> / <c>conflicts</c> / <c>reason</c> reached nothing it reads.</para>
/// </summary>
public sealed class GitIntegrateRunNode : INodeRuntime
{
Expand Down Expand Up @@ -79,15 +86,20 @@ public GitIntegrateRunNode(IBranchIntegrator integrator, IAgentWorkspaceResolver
"required": ["repositoryId"]
}
"""),
OutputSchema = SchemaBuilder.Parse("""
OutputSchema = SchemaBuilder.Parse($$"""
{
"type": "object",
"properties": {
"status": { "type": "string" },
"integratedBranch": { "type": ["string","null"] },
"appliedCount": { "type": "integer" },
"reason": { "type": ["string","null"] },
"conflicts": { "type": "array" }
"conflicts": { "type": "array" },
"summary": { "type": "string", "maxLength": {{RunIntegrationSummary.MaxChars}}, "description": "One account of what actually landed on the integrated branch, what conflicted or was skipped (and where its work is kept), and what was withheld — produced on every pass, and read by the plan-map synth." },
"withheld": { "type": "array", "items": { "type": "object", "properties": { "label": { "type": "string" }, "reason": { "type": "string" } } }, "description": "The units kept off the candidate BEFORE integration because their own acceptance check failed or was waived — in no other output. Empty when nothing was withheld." },
"reviewApproved": { "type": "boolean", "description": "Resumed pass only: the human's verdict on the parked conflict." },
"reviewComment": { "type": "string", "description": "Resumed pass only: the reviewer's comment." },
"reviewedBy": { "type": "string", "description": "Resumed pass only: who reviewed the parked conflict." }
}
}
""")
Expand All @@ -102,17 +114,18 @@ public async Task<NodeResult> RunAsync(NodeRunContext context, CancellationToken
var manifests = await _manifests.ListForWorkflowRunAsync(runId, teamId, cancellationToken).ConfigureAwait(false);
var agentWork = await LoadAgentWorkAsync(runId, teamId, cancellationToken).ConfigureAwait(false);
var contributions = RunIntegrationContributions.Build(repoId, manifests, agentWork);
var withheld = RunIntegrationContributions.Withheld(repoId, manifests, agentWork);

if (contributions.Count == 0)
return NodeResult.Ok(SkippedOutputs("the run produced no integrable work for this repository"));
return NodeResult.Ok(SkippedOutputs("the run produced no integrable work for this repository", withheld));

// The ancestor-most base, not the first contribution's: a withheld producer is dropped from the contributions
// while its manifest row survives, so the run's root lives in the ledger even when the surviving contributions
// are all dependents cut from a producer's head (see IntegrationBaseAnchor).
var baseSha = IntegrationBaseAnchor.Resolve(manifests, repoId, contributions.Select(c => c.BaseSha).FirstOrDefault(sha => !string.IsNullOrEmpty(sha)));

if (string.IsNullOrEmpty(baseSha))
return NodeResult.Ok(SkippedOutputs("the produced work recorded no base revision to integrate from"));
return NodeResult.Ok(SkippedOutputs("the produced work recorded no base revision to integrate from", withheld));

WorkspaceRequest? workspace;
try
Expand Down Expand Up @@ -165,7 +178,7 @@ public async Task<NodeResult> RunAsync(NodeRunContext context, CancellationToken

context.Logger.LogInformation("git.integrate_run on repo {RepoId}: {Status} ({Applied}/{Total} applied)", repoId, result.Status, result.AppliedCount, contributions.Count);

var outputs = GitIntegrateNode.ProjectOutputs(result);
var outputs = WithNarrative(GitIntegrateNode.ProjectOutputs(result), RunIntegrationSummary.ForResult(result, withheld), withheld);

// The review trail rides the outputs on the resumed pass — who looked, what they said — so the terminal
// (and any downstream consumer) sees the conflict was REVIEWED, never silently narrated past.
Expand Down Expand Up @@ -219,14 +232,28 @@ private async Task<IReadOnlyList<RunAgentWork>> LoadAgentWorkAsync(Guid runId, G
.Select(r => new RunAgentWork(r.Id, r.NodeId, r.IterationKey, r.CreatedDate, r.ResultJson, r.TaskJson))
.ToList();

private static Dictionary<string, JsonElement> SkippedOutputs(string reason) => new()
private static Dictionary<string, JsonElement> SkippedOutputs(string reason, IReadOnlyList<WithheldContribution> withheld)
{
["status"] = JsonSerializer.SerializeToElement("Skipped"),
["integratedBranch"] = JsonSerializer.SerializeToElement((string?)null),
["appliedCount"] = JsonSerializer.SerializeToElement(0),
["reason"] = JsonSerializer.SerializeToElement(reason),
["conflicts"] = JsonSerializer.SerializeToElement(Array.Empty<object>()),
};
var outputs = new Dictionary<string, JsonElement>
{
["status"] = JsonSerializer.SerializeToElement("Skipped"),
["integratedBranch"] = JsonSerializer.SerializeToElement((string?)null),
["appliedCount"] = JsonSerializer.SerializeToElement(0),
["reason"] = JsonSerializer.SerializeToElement(reason),
["conflicts"] = JsonSerializer.SerializeToElement(Array.Empty<object>()),
};

return WithNarrative(outputs, RunIntegrationSummary.ForSkipped(reason, withheld), withheld);
}

/// <summary>The outcome's prose and the units kept off the candidate, added to EVERY pass's outputs — the one seam, so no arm (clean, conflicted, skipped, resumed) can omit what the synth reads.</summary>
private static Dictionary<string, JsonElement> WithNarrative(Dictionary<string, JsonElement> outputs, string summary, IReadOnlyList<WithheldContribution> withheld)
{
outputs["summary"] = JsonSerializer.SerializeToElement(summary);
outputs["withheld"] = JsonSerializer.SerializeToElement(withheld.Select(w => new { label = w.Label, reason = w.Reason }));

return outputs;
}

private static bool TryReadGuid(NodeRunContext context, string key, out Guid id)
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -48,14 +48,7 @@ public static class RunIntegrationContributions
{
public static IReadOnlyList<BranchContribution> Build(Guid repositoryId, IReadOnlyList<PublishManifest> manifests, IReadOnlyList<RunAgentWork> agentWork)
{
var workByRunId = agentWork.ToDictionary(w => w.AgentRunId);

var produced = manifests
.Where(m => m.Kind == PublishManifestKind.Agent && m.AgentRunId is not null && m.RepositoryId == repositoryId && m.PublishStateValue != PublishState.None)
.Where(m => !IsWithheldFromHead(m))
.Select(m => (Manifest: m, Work: workByRunId.GetValueOrDefault(m.AgentRunId!.Value)))
.Where(pair => pair.Work is not null)
.Select(pair => (pair.Manifest, Work: pair.Work!));
var produced = ProducedRows(repositoryId, manifests, agentWork).Where(pair => !IsWithheldFromHead(pair.Manifest));

return LatestAttemptPerUnit(produced)
.OrderBy(pair => pair.Work.CreatedDate).ThenBy(pair => pair.Work.AgentRunId)
Expand All @@ -71,6 +64,56 @@ public static IReadOnlyList<BranchContribution> Build(Guid repositoryId, IReadOn
.ToList();
}

/// <summary>
/// The other face of <see cref="IsWithheldFromHead"/>: the units the head gate kept OUT of <see cref="Build"/>'s
/// contributions, reported instead of discarded. A withheld unit reaches no integrator outcome at all — it is
/// dropped before integration — so without this an outcome that reads <c>Clean</c> over the survivors is
/// indistinguishable from a run that integrated everything it produced, and whoever narrates it (the plan-map synth)
/// calls a partial deliverable whole.
///
/// <para>A unit is withheld only when NONE of its work landed. A retry respawns a fresh agent run, so a unit whose
/// first attempt flunked and whose second passed has both rows in the ledger — and its work is on the candidate;
/// naming the flunked attempt would tell the reader that delivered work was not. "Did it land" is asked on the same
/// unit the contribution reduction uses: the (node, iteration) cell where the cell IS the unit, the attempt itself
/// in the supervisor lane, where the cell is a turn shared by K concurrent deliverables and a peer that landed must
/// not hide the one that was withheld. One entry per unit, at its latest attempt's verdict, in agent-run creation
/// order — the same total order <see cref="Build"/> applies, so the report repeats across builds.</para>
/// </summary>
public static IReadOnlyList<WithheldContribution> Withheld(Guid repositoryId, IReadOnlyList<PublishManifest> manifests, IReadOnlyList<RunAgentWork> agentWork)
{
var produced = ProducedRows(repositoryId, manifests, agentWork).ToList();
var landed = LatestAttemptPerUnit(produced.Where(pair => !IsWithheldFromHead(pair.Manifest))).Select(pair => UnitKey(pair.Work)).ToHashSet();

return produced
.Where(pair => IsWithheldFromHead(pair.Manifest) && !landed.Contains(UnitKey(pair.Work)))
.GroupBy(pair => UnitKey(pair.Work))
.Select(LatestRowOfUnit)
.OrderBy(pair => pair.Work.CreatedDate).ThenBy(pair => pair.Work.AgentRunId)
.Select(pair => new WithheldContribution(UnitLabel(pair.Work), $"acceptance {pair.Manifest.AcceptanceState}"))
.ToList();
}

/// <summary>Every row this repository's agents PRODUCED — Agent-kind, in this repository, carrying something (a None-state row left no trace), and joined to the agent-run row that holds its result bytes. Before any verdict is applied: <see cref="Build"/> drops the withheld ones, <see cref="Withheld"/> reports them.</summary>
private static IEnumerable<(PublishManifest Manifest, RunAgentWork Work)> ProducedRows(Guid repositoryId, IReadOnlyList<PublishManifest> manifests, IReadOnlyList<RunAgentWork> agentWork)
{
var workByRunId = agentWork.ToDictionary(w => w.AgentRunId);

return manifests
.Where(m => m.Kind == PublishManifestKind.Agent && m.AgentRunId is not null && m.RepositoryId == repositoryId && m.PublishStateValue != PublishState.None)
.Select(m => (Manifest: m, Work: workByRunId.GetValueOrDefault(m.AgentRunId!.Value)))
.Where(pair => pair.Work is not null)
.Select(pair => (pair.Manifest, Work: pair.Work!));
}

/// <summary>The unit's newest row — by agent-run creation, then id, then alias — so the pick is total and repeats across builds (the order <see cref="KeepLatestAttempt"/> uses, plus the alias tie-break between one attempt's own sibling rows).</summary>
private static (PublishManifest Manifest, RunAgentWork Work) LatestRowOfUnit(IEnumerable<(PublishManifest Manifest, RunAgentWork Work)> unit) =>
unit.OrderBy(pair => pair.Work.CreatedDate).ThenBy(pair => pair.Work.AgentRunId).ThenBy(pair => pair.Manifest.RepositoryAlias, StringComparer.Ordinal).Last();

private static string UnitLabel(RunAgentWork work) => AgentAcceptanceContract.UnitId(work.NodeId, work.IterationKey ?? "");

/// <summary>What "this unit" means when asking whether it landed: its <see cref="UnitLabel"/> where the cell is the unit, the agent run itself where it is not (the fence <see cref="LatestAttemptPerUnit"/> applies).</summary>
private static string UnitKey(RunAgentWork work) => CellIsTheUnit(work.TaskJson) ? UnitLabel(work) : work.AgentRunId.ToString("N");

/// <summary>
/// "This unit's work is WITHHELD from the reviewable head" — the ledger-row analogue of the supervisor lane's
/// <c>SupervisorOutcome.IsWithheldFromHead</c> (its per-unit grade rejected it, or a human WAIVED its
Expand Down
Loading
Loading