feat: wire real AgentOS execution

This commit is contained in:
-Puter
2026-07-28 16:52:49 +05:30
parent 4d9f6da41b
commit 4edd456b5b
17 changed files with 1090 additions and 308 deletions

View File

@@ -1,4 +1,7 @@
import { defaultCodingKitV0 } from "@code/primitives/resolver";
import { WorkAttemptExecutionError } from "@code/primitives/execution-runtime";
import type { WorkAttemptExecutionErrorReason } from "@code/primitives/execution-runtime";
import { defaultCodingKitV0, resolveOutcome } from "@code/primitives/resolver";
import type { AttemptClassification } from "@code/primitives/resolver";
import { WorkflowManager } from "@convex-dev/workflow";
import type { WorkflowId } from "@convex-dev/workflow";
import { ConvexError, v } from "convex/values";
@@ -8,9 +11,68 @@ import type { Doc, Id } from "./_generated/dataModel";
import { internalMutation, internalQuery, mutation } from "./_generated/server";
import type { MutationCtx } from "./_generated/server";
import { requireProjectMember } from "./authz";
import { forgeForHost } from "./gitConnectionData";
export const workflow = new WorkflowManager(components.workflow);
// Lease window for a running real attempt. Long enough to outlast a single
// agent turn; the workflow's own retry/lease reconciliation covers crashes.
const LEASE_MS = 5 * 60_000;
// Runtime failure reason -> durable attempt classification. The Convex
// workflow stores classifications, not provider-specific reasons, so the
// resolver and Work lifecycle stay forge-agnostic. `retryable` from the
// runtime is carried alongside and consulted by the kit retry policy.
const FAILURE_REASON_CLASSIFICATION = {
Authentication: "PermanentFailure",
Cancelled: "Cancelled",
HarnessFailed: "RetryableFailure",
InvalidInput: "PermanentFailure",
ProviderUnavailable: "RetryableFailure",
RepositoryFailed: "PermanentFailure",
Timeout: "RetryableFailure",
} as const satisfies Record<
WorkAttemptExecutionErrorReason,
AttemptClassification
>;
export const classifyFailure = (
reason: WorkAttemptExecutionErrorReason
): AttemptClassification => FAILURE_REASON_CLASSIFICATION[reason];
// Normalize an error thrown from the agent action into the durable failure
// triple (reason/retryable/summary). WorkAttemptExecutionError is the
// classified runtime failure; anything else is a Convex/infrastructure error
// treated as a transient HarnessFailed so the workflow retry policy decides.
const toExecutionFailure = (
error: unknown
): {
message: string;
reason: WorkAttemptExecutionErrorReason;
retryable: boolean;
} =>
error instanceof WorkAttemptExecutionError
? {
message: error.message,
reason: error.reason,
retryable: error.retryable,
}
: {
message: error instanceof Error ? error.message : "Execution failed",
reason: "HarnessFailed",
retryable: true,
};
const failureReasonValues = v.union(
v.literal("Authentication"),
v.literal("Cancelled"),
v.literal("HarnessFailed"),
v.literal("InvalidInput"),
v.literal("ProviderUnavailable"),
v.literal("RepositoryFailed"),
v.literal("Timeout")
);
const resolveReadySlice = async (
ctx: MutationCtx,
work: Doc<"works">,
@@ -36,6 +98,34 @@ const resolveReadySlice = async (
return slice;
};
// Deployment prerequisite gate for real execution. Rejects before any
// workflow/sandbox starts when the project lacks an attached, forge-matched
// Git connection with a cloneable source URL and default branch. Returns the
// validated connection so the caller can pass it to the agent unchanged.
const validateProjectDeployment = async (
ctx: MutationCtx,
work: Doc<"works">
): Promise<Doc<"gitConnections">> => {
const project = await ctx.db.get(work.projectId);
if (!project) {
throw new ConvexError("Project not found");
}
if (!project.gitConnectionId) {
throw new ConvexError("Connect Git credentials to this project first");
}
const connection = await ctx.db.get(project.gitConnectionId);
if (!connection) {
throw new ConvexError("Git connection not found");
}
const expected = forgeForHost(project.sourceHost);
if (expected && connection.provider !== expected) {
throw new ConvexError(
`Git credential provider (${connection.provider}) does not match this project's forge (${expected})`
);
}
return connection;
};
export const execute = workflow
.define({ args: { attemptId: v.id("workAttempts") } })
.handler(async (step, args): Promise<void> => {
@@ -57,9 +147,12 @@ export const execute = workflow
},
});
} catch (error) {
const failure = toExecutionFailure(error);
await step.runMutation(internal.workExecutionWorkflow.failAttempt, {
attemptId: args.attemptId,
summary: error instanceof Error ? error.message : "Execution failed",
reason: failure.reason,
retryable: failure.retryable,
summary: failure.message,
});
}
});
@@ -91,10 +184,7 @@ export const startExecution = mutation({
) {
throw new ConvexError("Execution requires exact approved versions");
}
const project = await ctx.db.get(work.projectId);
if (!project?.gitConnectionId) {
throw new ConvexError("Connect Git credentials to this project first");
}
await validateProjectDeployment(ctx, work);
const slice = await resolveReadySlice(ctx, work, args.sliceId);
const createdAt = Date.now();
const runId = await ctx.db.insert("workRuns", {
@@ -159,6 +249,12 @@ export const executionContext = internalQuery({
if (!connection) {
throw new ConvexError("Git connection not found");
}
const expectedForge = forgeForHost(project.sourceHost);
if (expectedForge && connection.provider !== expectedForge) {
throw new ConvexError(
`Git credential provider (${connection.provider}) does not match this project's forge (${expectedForge})`
);
}
return {
attempt,
connection,
@@ -179,11 +275,16 @@ export const markAttemptRunning = internalMutation({
args: { attemptId: v.id("workAttempts") },
handler: async (ctx, args) => {
const attempt = await ctx.db.get(args.attemptId);
if (!attempt || attempt.status !== "queued") {
if (!attempt || attempt.status === "terminal") {
return false;
}
const run = await ctx.db.get(attempt.runId);
if (!run || run.status !== "running") {
return false;
}
await ctx.db.patch(attempt._id, {
startedAt: Date.now(),
leaseExpiresAt: Date.now() + LEASE_MS,
startedAt: attempt.startedAt ?? Date.now(),
status: "running",
});
return true;
@@ -193,7 +294,8 @@ export const markAttemptRunning = internalMutation({
const settleSliceAndWork = async (
ctx: MutationCtx,
run: Doc<"workRuns">,
succeeded: boolean
succeeded: boolean,
workStatus: Doc<"works">["status"]
) => {
if (!run.sliceRowId) {
return;
@@ -202,12 +304,18 @@ const settleSliceAndWork = async (
if (!slice) {
return;
}
await ctx.db.patch(slice._id, { status: succeeded ? "completed" : "ready" });
let sliceStatus: "blocked" | "completed" | "ready" = "ready";
if (succeeded) {
sliceStatus = "completed";
} else if (workStatus === "blocked") {
sliceStatus = "blocked";
}
await ctx.db.patch(slice._id, { status: sliceStatus });
const work = await ctx.db.get(run.workId);
if (!work) {
return;
}
let status: Doc<"works">["status"] = succeeded ? "completed" : "failed";
let status = workStatus;
if (succeeded) {
const slices = await ctx.db
.query("workSlices")
@@ -249,6 +357,8 @@ export const completeAttempt = internalMutation({
},
handler: async (ctx, args) => {
const attempt = await ctx.db.get(args.attemptId);
// Cancellation fencing: a terminal or cancelled attempt/run must never
// later settle success, even if the agent action raced to completion.
if (!attempt || attempt.status === "terminal") {
return;
}
@@ -257,6 +367,33 @@ export const completeAttempt = internalMutation({
if (!run || !work) {
throw new ConvexError("Execution records not found");
}
if (run.status === "terminal" || run.status === "cancelled") {
return;
}
// Empty-change rejection: a no-op result (no changed files, or identical
// base/candidate revisions) is not a successful implementation. Fail it
// as a permanent InvalidInput so the resolver does not mark Work done.
const noOp =
args.result.changedFiles.length === 0 ||
args.result.baseRevision === args.result.candidateRevision;
if (noOp) {
await ctx.db.patch(attempt._id, {
classification: "PermanentFailure",
endedAt: Date.now(),
failureReason: "InvalidInput",
leaseExpiresAt: undefined,
status: "terminal",
summary: "Execution produced no repository changes",
});
await ctx.db.patch(run._id, {
endedAt: Date.now(),
status: "terminal",
terminalClassification: "PermanentFailure",
terminalSummary: "Execution produced no repository changes",
});
await settleSliceAndWork(ctx, run, false, "failed");
return;
}
for (const item of args.result.events) {
await ctx.db.insert("workAttemptEvents", {
attemptId: attempt._id,
@@ -271,6 +408,7 @@ export const completeAttempt = internalMutation({
await ctx.db.patch(attempt._id, {
classification: "Succeeded",
endedAt,
leaseExpiresAt: undefined,
status: "terminal",
summary: args.result.summary,
});
@@ -309,7 +447,7 @@ export const completeAttempt = internalMutation({
}
: {}),
});
await settleSliceAndWork(ctx, run, true);
await settleSliceAndWork(ctx, run, true, "completed");
await ctx.db.insert("workEvents", {
createdAt: endedAt,
idempotencyKey: `real-run-completed:${run._id}`,
@@ -321,31 +459,107 @@ export const completeAttempt = internalMutation({
},
});
export const failAttempt = internalMutation({
args: { attemptId: v.id("workAttempts"), summary: v.string() },
export const launchAttemptWorkflow = internalMutation({
args: { attemptId: v.id("workAttempts") },
handler: async (ctx, args) => {
const attempt = await ctx.db.get(args.attemptId);
if (!attempt || attempt.status === "terminal") {
return;
}
const run = await ctx.db.get(attempt.runId);
if (!run) {
if (!run || run.status !== "running") {
return;
}
const workflowId: WorkflowId = await workflow.start(
ctx,
internal.workExecutionWorkflow.execute,
args
);
await ctx.db.patch(run._id, { workflowId });
},
});
export const failAttempt = internalMutation({
args: {
attemptId: v.id("workAttempts"),
reason: failureReasonValues,
retryable: v.boolean(),
summary: v.string(),
},
handler: async (ctx, args) => {
const attempt = await ctx.db.get(args.attemptId);
// Cancellation fencing: once an attempt or run is terminal/cancelled, a
// late failure from the workflow catch block must not overwrite the
// settled state. This also prevents a cancelled attempt from settling
// as a different classification after cancelExecution ran.
if (!attempt || attempt.status === "terminal") {
return;
}
const run = await ctx.db.get(attempt.runId);
if (!run || run.status === "terminal" || run.status === "cancelled") {
return;
}
const classification = classifyFailure(args.reason);
const endedAt = Date.now();
await ctx.db.patch(attempt._id, {
classification: "PermanentFailure",
classification,
endedAt,
failureReason: args.reason,
leaseExpiresAt: undefined,
status: "terminal",
summary: args.summary,
});
const outcome = resolveOutcome(
{ classification, retryable: args.retryable, summary: args.summary },
attempt.number,
defaultCodingKitV0.retryPolicy
);
await ctx.db.insert("resolverDecisions", {
attemptId: attempt._id,
attemptNumber: attempt.number,
classification,
createdAt: endedAt,
decision: outcome.kind,
...(outcome.kind === "terminal"
? { resultingWorkStatus: outcome.workStatus }
: {}),
runId: run._id,
summary: args.summary,
workId: attempt.workId,
});
if (outcome.kind === "retry") {
const nextAttemptId = await ctx.db.insert("workAttempts", {
number: attempt.number + 1,
runId: run._id,
status: "queued",
workId: attempt.workId,
workspaceKey: attempt.workspaceKey,
});
await ctx.scheduler.runAfter(
0,
internal.workExecutionWorkflow.launchAttemptWorkflow,
{ attemptId: nextAttemptId }
);
return;
}
await ctx.db.patch(run._id, {
endedAt,
status: "terminal",
terminalClassification: "PermanentFailure",
terminalClassification: classification,
terminalSummary: args.summary,
});
await settleSliceAndWork(ctx, run, false);
await settleSliceAndWork(ctx, run, false, outcome.workStatus);
const work = await ctx.db.get(run.workId);
if (work) {
await ctx.db.insert("workEvents", {
createdAt: endedAt,
idempotencyKey: `real-run-failed:${run._id}`,
kind: "run.completed",
payloadJson: JSON.stringify({ classification }),
referenceId: String(run._id),
workId: work._id,
});
}
},
});
@@ -375,15 +589,30 @@ export const cancelExecution = mutation({
0,
internal.workExecutionAgent.cancelAttempt,
{
attemptId: String(attempt._id),
workspaceKey: attempt.workspaceKey,
}
);
const cancelledAt = Date.now();
await ctx.db.patch(attempt._id, {
classification: "Cancelled",
endedAt: Date.now(),
endedAt: cancelledAt,
failureReason: "Cancelled",
leaseExpiresAt: undefined,
status: "terminal",
summary: "Execution cancelled",
});
await ctx.db.insert("resolverDecisions", {
attemptId: attempt._id,
attemptNumber: attempt.number,
classification: "Cancelled",
createdAt: cancelledAt,
decision: "terminal",
resultingWorkStatus: "ready",
runId: run._id,
summary: "Execution cancelled",
workId: run.workId,
});
}
await ctx.db.patch(run._id, {
endedAt: Date.now(),