Files
zopu-code/packages/backend/convex/workPlanning.ts
2026-07-28 14:50:45 +05:30

871 lines
27 KiB
TypeScript

import type { WorkEventKind } from "@code/primitives/work";
import {
decodeDefinition,
validateDefinition,
WorkQuestion,
} from "@code/primitives/work-definition";
import type { WorkDefinition } from "@code/primitives/work-definition";
import { validateDesignPacket } from "@code/primitives/work-design";
import { makeFunctionReference } from "convex/server";
import { ConvexError, v } from "convex/values";
import { Effect, Schema } from "effect";
import type { Doc, Id } from "./_generated/dataModel";
import {
env,
internalAction,
internalMutation,
internalQuery,
mutation,
query,
} from "./_generated/server";
import type { MutationCtx, QueryCtx } from "./_generated/server";
import { requireProjectMember } from "./authz";
const plannerRef = makeFunctionReference<
"action",
{ organizationId: Id<"organizations">; workId: Id<"works"> }
>("workPlanning:runPlanner");
const plannerContextRef = makeFunctionReference<
"query",
{ workId: Id<"works"> },
{ objective: string; status: string; title: string } | null
>("workPlanning:getPlannerContext");
const plannerFailureRef = makeFunctionReference<
"mutation",
{ workId: Id<"works">; error: string },
null
>("workPlanning:recordPlannerFailure");
const parsePayload = (payloadJson: string): unknown => {
try {
return JSON.parse(payloadJson) as unknown;
} catch {
throw new ConvexError("Proposal payload must be valid JSON");
}
};
const objectPayload = (payloadJson: string): Record<string, unknown> => {
const value = parsePayload(payloadJson);
if (typeof value !== "object" || value === null || Array.isArray(value)) {
throw new ConvexError("Proposal payload must be an object");
}
return value as Record<string, unknown>;
};
const requireWork = async (
ctx: QueryCtx | MutationCtx,
workId: Id<"works">
): Promise<Doc<"works">> => {
const work = await ctx.db.get(workId);
if (!work) {
throw new ConvexError("Work not found");
}
return work;
};
const appendEvent = async (
ctx: MutationCtx,
workId: Id<"works">,
kind: WorkEventKind,
idempotencyKey: string,
referenceId?: string,
payloadJson?: string
) => {
const existing = await ctx.db
.query("workEvents")
.withIndex("by_work_and_idempotencyKey", (q) =>
q.eq("workId", workId).eq("idempotencyKey", idempotencyKey)
)
.unique();
if (existing) {
return existing._id;
}
return await ctx.db.insert("workEvents", {
createdAt: Date.now(),
idempotencyKey,
kind,
workId,
...(referenceId ? { referenceId } : {}),
...(payloadJson ? { payloadJson } : {}),
});
};
const invalidateApprovalsAndDesign = async (
ctx: MutationCtx,
workId: Id<"works">
) => {
const approvals = await ctx.db
.query("workApprovals")
.withIndex("by_workId_and_kind", (q) => q.eq("workId", workId))
.collect();
for (const approval of approvals) {
if (approval.status === "active") {
await ctx.db.patch(approval._id, { status: "invalidated" });
}
}
const designs = await ctx.db
.query("designPackets")
.withIndex("by_work_and_version", (q) => q.eq("workId", workId))
.collect();
for (const design of designs) {
if (design.status === "current") {
await ctx.db.patch(design._id, { status: "superseded" });
}
}
};
export const requestDefinition = mutation({
args: { workId: v.id("works") },
handler: async (ctx, args) => {
const work = await requireWork(ctx, args.workId);
await requireProjectMember(ctx, work.projectId);
if (
work.status !== "proposed" &&
work.status !== "defining" &&
!(
work.status === "blocked" &&
work.definitionApprovalVersion !== work.definitionVersion
)
) {
throw new ConvexError("Work is not available for definition");
}
const now = Date.now();
await ctx.db.patch(work._id, { status: "defining", updatedAt: now });
await appendEvent(
ctx,
work._id,
"definition.requested",
`definition-requested:${work._id}`
);
await ctx.scheduler.runAfter(0, plannerRef, {
organizationId: work.organizationId,
workId: work._id,
});
return { status: "defining" as const };
},
});
const saveDefinition = async (
ctx: MutationCtx,
work: Doc<"works">,
payloadJson: string,
createdBy: string
) => {
if (work.status === "executing") {
throw new ConvexError("Cannot revise Work while a Run is executing");
}
const decoded = await Effect.runPromise(
decodeDefinition({
...objectPayload(payloadJson),
version: (work.definitionVersion ?? 0) + 1,
})
).catch((error: unknown) => {
throw new ConvexError(
error instanceof Error ? error.message : "Invalid Definition"
);
});
const { version } = decoded;
await invalidateApprovalsAndDesign(ctx, work._id);
const previous = await ctx.db
.query("workDefinitions")
.withIndex("by_work_and_version", (q) => q.eq("workId", work._id))
.collect();
for (const row of previous) {
if (row.status === "current") {
await ctx.db.patch(row._id, { status: "superseded" });
}
}
const definitionId = await ctx.db.insert("workDefinitions", {
createdAt: Date.now(),
createdBy,
payloadJson: JSON.stringify(decoded),
risk: decoded.risk,
status: "current",
version,
workId: work._id,
});
for (const question of decoded.questions) {
await ctx.db.insert("workQuestions", {
alternativesJson: JSON.stringify(question.alternatives),
answer: question.answer,
createdAt: Date.now(),
definitionVersion: version,
impact: question.impact,
prompt: question.prompt,
questionId: question.id,
recommendation: question.recommendation,
status: question.status,
workId: work._id,
});
}
await ctx.db.patch(work._id, {
definitionApprovalVersion: undefined,
definitionVersion: version,
designApprovalVersion: undefined,
designVersion: undefined,
status: "awaiting-definition-approval",
updatedAt: Date.now(),
});
await appendEvent(
ctx,
work._id,
work.definitionVersion ? "definition.revised" : "definition.saved",
`definition:${version}`,
String(definitionId),
JSON.stringify(decoded)
);
return { definitionId, version };
};
export const saveDefinitionProposal = mutation({
args: { payloadJson: v.string(), workId: v.id("works") },
handler: async (ctx, args) => {
const work = await requireWork(ctx, args.workId);
await requireProjectMember(ctx, work.projectId);
return await saveDefinition(ctx, work, args.payloadJson, "user");
},
});
export const reviseDefinition = saveDefinitionProposal;
export const approveDefinition = mutation({
args: { version: v.number(), workId: v.id("works") },
handler: async (ctx, args) => {
const work = await requireWork(ctx, args.workId);
await requireProjectMember(ctx, work.projectId);
if (
work.definitionVersion !== args.version ||
work.status !== "awaiting-definition-approval"
) {
throw new ConvexError("Definition version is not current or approvable");
}
const row = await ctx.db
.query("workDefinitions")
.withIndex("by_work_and_version", (q) =>
q.eq("workId", work._id).eq("version", args.version)
)
.unique();
if (!row) {
throw new ConvexError("Definition not found");
}
const definition = JSON.parse(row.payloadJson) as WorkDefinition;
const valid = await Effect.runPromise(validateDefinition(definition)).catch(
(error: unknown) => {
throw new ConvexError(
error instanceof Error
? error.message
: "Definition cannot be approved"
);
}
);
const questions = await ctx.db
.query("workQuestions")
.withIndex("by_workId_and_definitionVersion", (q) =>
q.eq("workId", work._id).eq("definitionVersion", args.version)
)
.collect();
if (
questions.some(
(question) => question.status === "open" && question.impact === "high"
)
) {
throw new ConvexError(
"High-impact open questions must be resolved before approval"
);
}
const identity = await ctx.auth.getUserIdentity();
if (!identity) {
throw new ConvexError("Authentication required");
}
await ctx.db.insert("workApprovals", {
approvedAt: Date.now(),
approvedBy: identity.tokenIdentifier,
definitionVersion: args.version,
kind: "definition",
status: "active",
workId: work._id,
});
await ctx.db.patch(work._id, {
definitionApprovalVersion: valid.version,
status: "designing",
updatedAt: Date.now(),
});
await appendEvent(
ctx,
work._id,
"definition.approved",
`definition-approved:${args.version}`
);
await ctx.scheduler.runAfter(0, plannerRef, {
organizationId: work.organizationId,
workId: work._id,
});
return { status: "designing" as const, version: args.version };
},
});
const reviseQuestion = async (
ctx: MutationCtx,
work: Doc<"works">,
questionId: string,
update: { status: "answered" | "withdrawn"; answer?: string }
) => {
const current = await ctx.db
.query("workDefinitions")
.withIndex("by_work_and_version", (q) =>
q.eq("workId", work._id).eq("version", work.definitionVersion ?? 0)
)
.unique();
if (!current) {
throw new ConvexError("Current Definition not found");
}
const definition = JSON.parse(current.payloadJson) as WorkDefinition;
const persistedQuestions = await ctx.db
.query("workQuestions")
.withIndex("by_workId_and_definitionVersion", (q) =>
q.eq("workId", work._id).eq("definitionVersion", current.version)
)
.collect();
if (
!persistedQuestions.some((question) => question.questionId === questionId)
) {
throw new ConvexError("Question not found");
}
const next = {
...definition,
questions: persistedQuestions.map((question) => ({
alternatives: JSON.parse(question.alternativesJson) as string[],
answer:
question.questionId === questionId ? update.answer : question.answer,
id: question.questionId,
impact: question.impact,
prompt: question.prompt,
recommendation: question.recommendation,
status:
question.questionId === questionId ? update.status : question.status,
})),
};
return await saveDefinition(ctx, work, JSON.stringify(next), "user");
};
export const answerQuestion = mutation({
args: { answer: v.string(), questionId: v.string(), workId: v.id("works") },
handler: async (ctx, args) => {
const work = await requireWork(ctx, args.workId);
await requireProjectMember(ctx, work.projectId);
const result = await reviseQuestion(ctx, work, args.questionId, {
answer: args.answer,
status: "answered",
});
await appendEvent(
ctx,
work._id,
"question.answered",
`question-answered:${args.questionId}:${result.version}`
);
return result;
},
});
export const withdrawQuestion = mutation({
args: { questionId: v.string(), workId: v.id("works") },
handler: async (ctx, args) => {
const work = await requireWork(ctx, args.workId);
await requireProjectMember(ctx, work.projectId);
const result = await reviseQuestion(ctx, work, args.questionId, {
status: "withdrawn",
});
await appendEvent(
ctx,
work._id,
"question.withdrawn",
`question-withdrawn:${args.questionId}:${result.version}`
);
return result;
},
});
export const requestDesign = mutation({
args: { workId: v.id("works") },
handler: async (ctx, args) => {
const work = await requireWork(ctx, args.workId);
await requireProjectMember(ctx, work.projectId);
if (
(work.status !== "designing" && work.status !== "blocked") ||
work.definitionApprovalVersion !== work.definitionVersion
) {
throw new ConvexError("Definition must be approved before Design");
}
if (work.status === "blocked") {
await ctx.db.patch(work._id, {
status: "designing",
updatedAt: Date.now(),
});
}
await appendEvent(
ctx,
work._id,
"design.requested",
`design-requested:${work.definitionVersion}`
);
await ctx.scheduler.runAfter(0, plannerRef, {
organizationId: work.organizationId,
workId: work._id,
});
return { status: "designing" as const };
},
});
const saveDesign = async (
ctx: MutationCtx,
work: Doc<"works">,
payloadJson: string,
createdBy: string
) => {
if (work.status === "executing") {
throw new ConvexError("Cannot revise Work while a Run is executing");
}
if (work.definitionApprovalVersion !== work.definitionVersion) {
throw new ConvexError("Design must bind an approved Definition version");
}
const decoded = await Effect.runPromise(
validateDesignPacket({
...objectPayload(payloadJson),
definitionVersion: work.definitionVersion,
version: (work.designVersion ?? 0) + 1,
})
).catch((error: unknown) => {
throw new ConvexError(
error instanceof Error ? error.message : "Invalid Design Packet"
);
});
const prior = await ctx.db
.query("designPackets")
.withIndex("by_work_and_version", (q) => q.eq("workId", work._id))
.collect();
for (const row of prior) {
if (row.status === "current") {
await ctx.db.patch(row._id, { status: "superseded" });
}
}
const designApprovals = await ctx.db
.query("workApprovals")
.withIndex("by_workId_and_kind", (q) =>
q.eq("workId", work._id).eq("kind", "design")
)
.collect();
for (const approval of designApprovals) {
if (approval.status === "active") {
await ctx.db.patch(approval._id, { status: "invalidated" });
}
}
const designId = await ctx.db.insert("designPackets", {
createdAt: Date.now(),
createdBy,
definitionVersion: decoded.definitionVersion,
payloadJson: JSON.stringify(decoded),
status: "current",
version: decoded.version,
workId: work._id,
});
for (const [ordinal, slice] of decoded.slices.entries()) {
await ctx.db.insert("workSlices", {
createdAt: Date.now(),
designVersion: decoded.version,
objective: slice.objective,
observableBehavior: slice.observableBehavior,
ordinal,
payloadJson: JSON.stringify(slice),
sliceId: slice.id,
status: ordinal === 0 ? "ready" : "planned",
title: slice.title,
workId: work._id,
});
}
await ctx.db.patch(work._id, {
designApprovalVersion: undefined,
designVersion: decoded.version,
status: "awaiting-design-approval",
updatedAt: Date.now(),
});
await appendEvent(
ctx,
work._id,
work.designVersion ? "design.revised" : "design.saved",
`design:${decoded.version}`,
String(designId),
JSON.stringify(decoded)
);
return { designId, version: decoded.version };
};
export const saveDesignProposal = mutation({
args: { payloadJson: v.string(), workId: v.id("works") },
handler: async (ctx, args) => {
const work = await requireWork(ctx, args.workId);
await requireProjectMember(ctx, work.projectId);
return await saveDesign(ctx, work, args.payloadJson, "user");
},
});
export const reviseDesign = saveDesignProposal;
export const approveDesign = mutation({
args: {
definitionVersion: v.number(),
designVersion: v.number(),
workId: v.id("works"),
},
handler: async (ctx, args) => {
const work = await requireWork(ctx, args.workId);
await requireProjectMember(ctx, work.projectId);
if (
work.definitionApprovalVersion !== args.definitionVersion ||
work.designVersion !== args.designVersion ||
work.status !== "awaiting-design-approval"
) {
throw new ConvexError(
"Design approval must bind current Definition and Design versions"
);
}
const identity = await ctx.auth.getUserIdentity();
if (!identity) {
throw new ConvexError("Authentication required");
}
await ctx.db.insert("workApprovals", {
approvedAt: Date.now(),
approvedBy: identity.tokenIdentifier,
definitionVersion: args.definitionVersion,
designVersion: args.designVersion,
kind: "design",
status: "active",
workId: work._id,
});
await ctx.db.patch(work._id, {
designApprovalVersion: args.designVersion,
status: "ready",
updatedAt: Date.now(),
});
await appendEvent(
ctx,
work._id,
"design.approved",
`design-approved:${args.definitionVersion}:${args.designVersion}`
);
return { status: "ready" as const };
},
});
export const submitDefinitionProposal = mutation({
args: { payloadJson: v.string(), token: v.string(), workId: v.id("works") },
handler: async (ctx, args) => {
if (args.token !== env.FLUE_DB_TOKEN) {
throw new ConvexError("Invalid agent control token");
}
const work = await requireWork(ctx, args.workId);
return await saveDefinition(ctx, work, args.payloadJson, "work-planner");
},
});
export const submitDesignProposal = mutation({
args: { payloadJson: v.string(), token: v.string(), workId: v.id("works") },
handler: async (ctx, args) => {
if (args.token !== env.FLUE_DB_TOKEN) {
throw new ConvexError("Invalid agent control token");
}
const work = await requireWork(ctx, args.workId);
return await saveDesign(ctx, work, args.payloadJson, "work-planner");
},
});
export const submitQuestion = mutation({
args: { questionJson: v.string(), token: v.string(), workId: v.id("works") },
handler: async (ctx, args) => {
if (args.token !== env.FLUE_DB_TOKEN) {
throw new ConvexError("Invalid agent control token");
}
const work = await requireWork(ctx, args.workId);
if (work.definitionVersion === undefined) {
throw new ConvexError("Work has no Definition to attach a question to");
}
if (
work.status !== "defining" &&
work.status !== "awaiting-definition-approval"
) {
throw new ConvexError(
"Questions can only be submitted while the Work is being defined"
);
}
const question = await Effect.runPromise(
Schema.decodeUnknownEffect(WorkQuestion)(parsePayload(args.questionJson))
).catch((error: unknown) => {
throw new ConvexError(
error instanceof Error ? error.message : "Invalid question"
);
});
const existing = await ctx.db
.query("workQuestions")
.withIndex("by_workId_and_definitionVersion_and_questionId", (q) =>
q
.eq("workId", work._id)
.eq("definitionVersion", work.definitionVersion!)
.eq("questionId", question.id)
)
.unique();
const questionRow = {
alternativesJson: JSON.stringify(question.alternatives),
answer: question.answer,
createdAt: Date.now(),
definitionVersion: work.definitionVersion,
impact: question.impact,
prompt: question.prompt,
questionId: question.id,
recommendation: question.recommendation,
status: question.status,
workId: work._id,
};
if (existing) {
const same =
existing.prompt === questionRow.prompt &&
existing.impact === questionRow.impact &&
existing.recommendation === questionRow.recommendation &&
existing.alternativesJson === questionRow.alternativesJson &&
existing.status === questionRow.status &&
existing.answer === questionRow.answer;
if (!same) {
throw new ConvexError("Question ID already has different content");
}
return { accepted: false };
}
await ctx.db.insert("workQuestions", questionRow);
await appendEvent(
ctx,
work._id,
"question.created",
`question-created:${work.definitionVersion}:${question.id}`,
question.id,
JSON.stringify(question)
);
return { accepted: true };
},
});
export const recordPlannerFailure = internalMutation({
args: { error: v.string(), workId: v.id("works") },
handler: async (ctx, args): Promise<null> => {
const work = await ctx.db.get(args.workId);
if (!work || (work.status !== "defining" && work.status !== "designing")) {
return null;
}
await ctx.db.patch(work._id, { status: "blocked", updatedAt: Date.now() });
await appendEvent(
ctx,
work._id,
"planner.failed",
`planner-failed:${work._id}:${work.status}`,
undefined,
JSON.stringify({ error: args.error, phase: work.status })
);
return null;
},
});
export const runPlanner = internalAction({
args: { organizationId: v.id("organizations"), workId: v.id("works") },
handler: async (ctx, args) => {
const flueUrl = (env as typeof env & { readonly FLUE_URL?: string })
.FLUE_URL;
if (!flueUrl) {
await ctx.runMutation(plannerFailureRef, {
error: "FLUE_URL is not configured in Convex",
workId: args.workId,
});
return null;
}
const context = await ctx.runQuery(plannerContextRef, {
workId: args.workId,
});
if (!context) {
return null;
}
const endpoint = new URL(
`agents/work-planner/${encodeURIComponent(String(args.organizationId))}`,
`${flueUrl.replace(/\/+$/u, "")}/`
);
endpoint.searchParams.set("wait", "result");
try {
const response = await fetch(endpoint, {
body: JSON.stringify({
message: `Plan Work ${String(args.workId)} titled "${context.title}" with objective "${context.objective}". Current state is ${context.status}. If the Work is defining, submit only a Definition proposal and questions. If the Work is designing, submit only a Design Packet proposal bound to the approved Definition. Never approve or execute.`,
}),
headers: {
authorization: `Bearer ${env.FLUE_DB_TOKEN}`,
"content-type": "application/json",
"x-zopu-organization-id": String(args.organizationId),
"x-zopu-request-id": `work-planner:${args.workId}`,
},
method: "POST",
});
if (response.ok) {
return null;
}
await ctx.runMutation(plannerFailureRef, {
error: `Work Planner request failed (${response.status})`,
workId: args.workId,
});
} catch (error) {
await ctx.runMutation(plannerFailureRef, {
error: error instanceof Error ? error.message : String(error),
workId: args.workId,
});
}
// The private worker submits typed proposals through Convex mutations; it
// never approves or advances Work itself.
return null;
},
});
export const getPlannerContext = internalQuery({
args: { workId: v.id("works") },
handler: async (ctx, args) => {
const work = await ctx.db.get(args.workId);
return work
? { objective: work.objective, status: work.status, title: work.title }
: null;
},
});
export const listForProject = query({
args: { projectId: v.id("projects") },
handler: async (ctx, args) => {
await requireProjectMember(ctx, args.projectId);
const works = await ctx.db
.query("works")
.withIndex("by_project_and_createdAt", (q: any) =>
q.eq("projectId", args.projectId)
)
.order("desc")
.take(100);
return await Promise.all(
works.map(async (work) => {
const definitions = await ctx.db
.query("workDefinitions")
.withIndex("by_work_and_version", (q: any) =>
q.eq("workId", work._id)
)
.order("desc")
.take(10);
const designs = await ctx.db
.query("designPackets")
.withIndex("by_work_and_version", (q: any) =>
q.eq("workId", work._id)
)
.order("desc")
.take(10);
const slices = await ctx.db
.query("workSlices")
.withIndex("by_workId_and_designVersion", (q) =>
q
.eq("workId", work._id)
.eq("designVersion", work.designVersion ?? 0)
)
.collect();
const runRows = await ctx.db
.query("workRuns")
.withIndex("by_work_and_createdAt", (q: any) =>
q.eq("workId", work._id)
)
.order("desc")
.take(10);
const runs = await Promise.all(
runRows.map(async (run) => {
const attempts = await ctx.db
.query("workAttempts")
.withIndex("by_runId_and_number", (q) => q.eq("runId", run._id))
.collect();
const eventsByAttempt = await Promise.all(
attempts.map((attempt) =>
ctx.db
.query("workAttemptEvents")
.withIndex("by_attempt_and_sequence", (q: any) =>
q.eq("attemptId", attempt._id)
)
.collect()
)
);
const attemptEvents = eventsByAttempt
.flat()
.sort((left, right) => left.occurredAt - right.occurredAt);
const artifacts = await ctx.db
.query("workArtifacts")
.withIndex("by_runId_and_createdAt", (q) =>
q.eq("runId", run._id)
)
.collect();
return { ...run, artifacts, attemptEvents, attempts };
})
);
const events = await ctx.db
.query("workEvents")
.withIndex("by_work_and_createdAt", (q: any) =>
q.eq("workId", work._id)
)
.order("desc")
.take(100);
const attachments = await ctx.db
.query("signalWorkAttachments")
.withIndex("by_work", (q: any) => q.eq("workId", work._id))
.collect();
const signals = await Promise.all(
attachments.map(async (attachment) => {
const signal = await ctx.db.get(attachment.signalId);
if (!signal) {
return null;
}
const sources = await ctx.db
.query("signalSources")
.withIndex("by_signalId_and_ordinal", (q) =>
q.eq("signalId", signal._id)
)
.collect();
return {
createdAt: signal.createdAt,
signalId: signal._id,
sources: sources
.sort((a, b) => a.ordinal - b.ordinal)
.map((source) => ({
createdAt: source.sourceCreatedAt,
messageId: String(source.messageId),
rawText: source.rawTextSnapshot,
submissionId: null,
})),
summary: signal.summary,
title: signal.title,
};
})
);
const currentDefinition = definitions.find(
(definition) => definition.status === "current"
);
const currentDesign = designs.find(
(design) => design.status === "current"
);
return {
...work,
definition: currentDefinition
? JSON.parse(currentDefinition.payloadJson)
: null,
definitions,
design: currentDesign ? JSON.parse(currentDesign.payloadJson) : null,
designs,
events,
runs,
signals: signals.filter((signal) => signal !== null),
slices: slices.sort((a, b) => a.ordinal - b.ordinal),
};
})
);
},
});