Files
zopu-code/packages/backend/convex/works.ts

271 lines
7.8 KiB
TypeScript

import { env } from "@code/env/convex";
import {
signalAttachedEvent,
workDraftFromSignal,
workProposedEventFromSignal,
} from "@code/primitives/work";
import { ConvexError, v } from "convex/values";
import { Effect } from "effect";
import type { Doc, Id } from "./_generated/dataModel";
import { mutation, type MutationCtx, query } from "./_generated/server";
import { requireProjectMember } from "./authz";
const requireAgent = (token: string) => {
if (token !== env.FLUE_DB_TOKEN) {
throw new ConvexError("Invalid agent control token");
}
};
const requireSignal = async (
ctx: MutationCtx,
organizationId: Id<"organizations">,
signalId: Id<"signals">
): Promise<Doc<"signals">> => {
const signal = await ctx.db.get(signalId);
if (!signal || signal.organizationId !== organizationId) {
throw new ConvexError("Signal not found");
}
if (!signal.projectId) {
throw new ConvexError("Signal must belong to a project");
}
return signal;
};
const requireWork = async (
ctx: MutationCtx,
organizationId: Id<"organizations">,
workId: Id<"works">
): Promise<Doc<"works">> => {
const work = await ctx.db.get(workId);
if (!work || work.organizationId !== organizationId) {
throw new ConvexError("Work not found");
}
return work;
};
const attachSignal = async (
ctx: MutationCtx,
signal: Doc<"signals">,
work: Doc<"works">
): Promise<boolean> => {
const event = await Effect.runPromise(
signalAttachedEvent({
signal: {
organizationId: String(signal.organizationId),
projectId: String(signal.projectId),
signalId: String(signal._id),
title: signal.title,
},
work: {
organizationId: String(work.organizationId),
projectId: String(work.projectId),
workId: String(work._id),
},
})
).catch((error: unknown) => {
throw new ConvexError(
error instanceof Error ? error.message : "Invalid Work attachment"
);
});
const existing = await ctx.db
.query("signalWorkAttachments")
.withIndex("by_signal_and_work", (q) =>
q.eq("signalId", signal._id).eq("workId", work._id)
)
.unique();
if (existing) {
return false;
}
const createdAt = Date.now();
await ctx.db.insert("signalWorkAttachments", {
createdAt,
signalId: signal._id,
workId: work._id,
});
await ctx.db.insert("workEvents", {
createdAt,
idempotencyKey: event.idempotencyKey,
kind: event.kind,
signalId: signal._id,
workId: work._id,
});
await ctx.db.patch(work._id, { updatedAt: createdAt });
return true;
};
export const listProposedForAgent = query({
args: {
organizationId: v.id("organizations"),
projectId: v.id("projects"),
token: v.string(),
},
handler: async (ctx, args) => {
requireAgent(args.token);
const project = await ctx.db.get(args.projectId);
if (!project || project.organizationId !== args.organizationId) {
throw new ConvexError("Project not found");
}
return await ctx.db
.query("works")
.withIndex("by_project_and_createdAt", (q) =>
q.eq("projectId", args.projectId)
)
.order("desc")
.take(50);
},
});
export const createFromSignal = mutation({
args: {
organizationId: v.id("organizations"),
signalId: v.id("signals"),
token: v.string(),
},
handler: async (ctx, args) => {
requireAgent(args.token);
const signal = await requireSignal(ctx, args.organizationId, args.signalId);
const existingAttachment = await ctx.db
.query("signalWorkAttachments")
.withIndex("by_signal", (q) => q.eq("signalId", signal._id))
.first();
if (existingAttachment) {
return { created: false, workId: existingAttachment.workId };
}
const constraints = await ctx.db
.query("signalConstraints")
.withIndex("by_signalId_and_ordinal", (q) => q.eq("signalId", signal._id))
.collect();
const draft = await Effect.runPromise(
workDraftFromSignal({
constraints: constraints
.sort((left, right) => left.ordinal - right.ordinal)
.map((constraint) => constraint.value),
desiredOutcome: signal.desiredOutcome,
summary: signal.summary,
title: signal.title,
})
).catch((error: unknown) => {
throw new ConvexError(
error instanceof Error ? error.message : "Invalid Work proposal"
);
});
const createdAt = Date.now();
const event = await Effect.runPromise(
workProposedEventFromSignal({
organizationId: String(signal.organizationId),
projectId: String(signal.projectId),
signalId: String(signal._id),
title: draft.title,
})
).catch((error: unknown) => {
throw new ConvexError(
error instanceof Error ? error.message : "Invalid Work event"
);
});
const workId = await ctx.db.insert("works", {
createdAt,
objective: draft.objective,
organizationId: signal.organizationId,
projectId: signal.projectId as Id<"projects">,
status: "proposed",
title: draft.title,
updatedAt: createdAt,
});
await ctx.db.insert("signalWorkAttachments", {
createdAt,
signalId: signal._id,
workId,
});
await ctx.db.insert("workEvents", {
createdAt,
idempotencyKey: event.idempotencyKey,
kind: event.kind,
signalId: signal._id,
workId,
});
return { created: true, workId };
},
});
export const attachSignalToWork = mutation({
args: {
organizationId: v.id("organizations"),
signalId: v.id("signals"),
token: v.string(),
workId: v.id("works"),
},
handler: async (ctx, args) => {
requireAgent(args.token);
const signal = await requireSignal(ctx, args.organizationId, args.signalId);
const work = await requireWork(ctx, args.organizationId, args.workId);
return {
attached: await attachSignal(ctx, signal, work),
workId: work._id,
};
},
});
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) =>
q.eq("projectId", args.projectId)
)
.order("desc")
.take(100);
return await Promise.all(
works.map(async (work) => {
const attachments = await ctx.db
.query("signalWorkAttachments")
.withIndex("by_work", (q) => 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,
summary: signal.summary,
sources: sources
.sort((a, b) => a.ordinal - b.ordinal)
.map((source) => ({
createdAt: source.sourceCreatedAt,
messageId: String(source.messageId),
rawText: source.rawTextSnapshot,
submissionId: null,
})),
title: signal.title,
};
})
);
const events = await ctx.db
.query("workEvents")
.withIndex("by_work_and_createdAt", (q) => q.eq("workId", work._id))
.order("desc")
.collect();
return {
...work,
events,
signals: signals.filter((signal) => signal !== null),
};
})
);
},
});