271 lines
7.8 KiB
TypeScript
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),
|
|
};
|
|
})
|
|
);
|
|
},
|
|
});
|