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 => { 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; }; const requireWork = async ( ctx: QueryCtx | MutationCtx, workId: Id<"works"> ): Promise> => { 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 => { 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), }; }) ); }, });