feat: dispatch authenticated project issue requests

This commit is contained in:
-Puter
2026-07-24 02:00:02 +05:30
parent 7974bd5f3e
commit 38c73847e0
11 changed files with 806 additions and 65 deletions

View File

@@ -149,10 +149,12 @@
"dependencies": {
"@code/backend": "workspace:*",
"@code/env": "workspace:*",
"@code/primitives": "workspace:*",
"@flue/runtime": "latest",
"@rivet-dev/agentos-core": "catalog:",
"convex": "catalog:",
"hono": "4.12.30",
"effect": "catalog:",
"hono": "4.12.31",
"valibot": "^1.4.2",
},
"devDependencies": {
@@ -2466,7 +2468,7 @@
"hoist-non-react-statics": ["hoist-non-react-statics@3.3.2", "", { "dependencies": { "react-is": "^16.7.0" } }, "sha512-/gGivxi8JPKWNm/W0jSmzcMPpfpPLc3dY/6GxhX2hQ9iGj3aDfklV4ET7NjKpSinLpJ5vafa9iiGIEZg10SfBw=="],
"hono": ["hono@4.12.30", "", {}, "sha512-emn+JoJjrN9YTpRDS5it/UI2SO9BAE37T6I3d963RxcZ81G9A4pr2SZTEiiaiKbzx+NKRg5BZ89fCL7gCJCUog=="],
"hono": ["hono@4.12.31", "", {}, "sha512-zJIHFrl6bq3RDd2YusFNCDlM8qUprxKswyi/OPzPyzKDdyBXDqWx8bZlZ7R+saTdSTatUmb3O7K4SspGPaEOQg=="],
"hono-openapi": ["hono-openapi@1.3.1", "", { "peerDependencies": { "@hono/standard-validator": "^0.2.0", "@standard-community/standard-json": "^0.3.5", "@standard-community/standard-openapi": "^0.2.9", "@types/json-schema": "^7.0.15", "hono": "^4.11.2", "openapi-types": "^12.1.3" }, "optionalPeers": ["@hono/standard-validator", "hono"] }, "sha512-NLVeVkhKZ3drmQNEIPac8HX8Y54uf1hJAgIM/7MfDsaeVVmB+QILWQxx5x3R3NvRHgedcbEbOCGY2uR7WQYyMw=="],
@@ -3854,8 +3856,6 @@
"@expo/xcpretty/chalk": ["chalk@4.1.2", "", { "dependencies": { "ansi-styles": "^4.1.0", "supports-color": "^7.1.0" } }, "sha512-oKnbhFyRIXpUuez8iBMmyEa4nbj4IOQyuhc/wy9kY7/WVPcwIO9VA668Pu8RkO7+0G76SLROeyw9CpQ061i4mA=="],
"@flue/runtime/hono": ["hono@4.12.31", "", {}, "sha512-zJIHFrl6bq3RDd2YusFNCDlM8qUprxKswyi/OPzPyzKDdyBXDqWx8bZlZ7R+saTdSTatUmb3O7K4SspGPaEOQg=="],
"@google/genai/google-auth-library": ["google-auth-library@10.9.0", "", { "dependencies": { "base64-js": "^1.3.0", "ecdsa-sig-formatter": "^1.0.11", "gaxios": "^7.1.4", "gcp-metadata": "8.1.2", "google-logging-utils": "1.1.3", "jws": "^4.0.0" } }, "sha512-xtvUqvINPhTaBm7nXqlYPcrMHJPm1lCNdSovxnKKhTm+4JsvQ+KGVYJViLoH9Yxu8w+T0Qv5HubzYT9BLrppJg=="],
"@google/genai/p-retry": ["p-retry@4.6.2", "", { "dependencies": { "@types/retry": "0.12.0", "retry": "^0.13.1" } }, "sha512-312Id396EbJdvRONlngUx0NydfrIQ5lsYu0znKVUzVvArzEIt08V1qhtyESbGVd1FGX7UKtiFp5uwKZdM8wIuQ=="],
@@ -3888,8 +3888,6 @@
"@modelcontextprotocol/sdk/@hono/node-server": ["@hono/node-server@1.19.14", "", { "peerDependencies": { "hono": "^4" } }, "sha512-GwtvgtXxnWsucXvbQXkRgqksiH2Qed37H9xHZocE5sA3N8O8O8/8FA3uclQXxXVzc9XBZuEOMK7+r02FmSpHtw=="],
"@modelcontextprotocol/sdk/hono": ["hono@4.12.31", "", {}, "sha512-zJIHFrl6bq3RDd2YusFNCDlM8qUprxKswyi/OPzPyzKDdyBXDqWx8bZlZ7R+saTdSTatUmb3O7K4SspGPaEOQg=="],
"@opentui/core/diff": ["diff@9.0.0", "", {}, "sha512-svtcdpS8CgJyqAjEQIXdb3OjhFVVYjzGAPO8WGCmRbrml64SPw/jJD4GoE98aR7r25A0XcgrK3F02yw9R/vhQw=="],
"@react-native/dev-middleware/open": ["open@7.4.2", "", { "dependencies": { "is-docker": "^2.0.0", "is-wsl": "^2.1.1" } }, "sha512-MVHddDVweXZF3awtlAS+6pgKLlm/JgxZ90+/NBurBoQctVOOB/zDdVjcyPzQ+0laDGbsWgrRkflI65sQeOgT9Q=="],
@@ -4182,8 +4180,6 @@
"ripemd160/hash-base": ["hash-base@3.1.2", "", { "dependencies": { "inherits": "^2.0.4", "readable-stream": "^2.3.8", "safe-buffer": "^5.2.1", "to-buffer": "^1.2.1" } }, "sha512-Bb33KbowVTIj5s7Ked1OsqHUeCpz//tPwR+E2zJgJKo9Z5XolZ9b6bdUgjmYlwnWhoOQKoTd1TYToZGn5mAYOg=="],
"rivetkit/hono": ["hono@4.12.31", "", {}, "sha512-zJIHFrl6bq3RDd2YusFNCDlM8qUprxKswyi/OPzPyzKDdyBXDqWx8bZlZ7R+saTdSTatUmb3O7K4SspGPaEOQg=="],
"rivetkit/uuid": ["uuid@12.0.1", "", { "bin": { "uuid": "dist/bin/uuid" } }, "sha512-9obBF8sMIHJWNQaO6IGOG8giGa/jUpKX34bz6o4whVs8M0WAvhID2tNxYp6A2XEBJPuZSX8wsS/6TEKfIDc+nw=="],
"router/path-to-regexp": ["path-to-regexp@8.4.2", "", {}, "sha512-qRcuIdP69NPm4qbACK+aDogI5CBDMi1jKe0ry5rSQJz8JVLsC7jV8XpiJjGRLLol3N+R5ihGYcrPLTno6pAdBA=="],
@@ -4376,8 +4372,6 @@
"@rivet-dev/agentos/rivetkit/@rivetkit/workflow-engine": ["@rivetkit/workflow-engine@2.3.7", "", { "dependencies": { "@rivetkit/bare-ts": "^0.6.2", "cbor-x": "^1.6.0", "fdb-tuple": "^1.0.0", "pino": "^9.6.0", "vbare": "^0.0.4" } }, "sha512-+C4kGuSNysrw2zs/LQoYxouOxufKHHMxiIXKF8Ab9GU2BCJ32RUST9dYU/gaQE/w9uM+hqAHUz0DTlOBYU6p6g=="],
"@rivet-dev/agentos/rivetkit/hono": ["hono@4.12.31", "", {}, "sha512-zJIHFrl6bq3RDd2YusFNCDlM8qUprxKswyi/OPzPyzKDdyBXDqWx8bZlZ7R+saTdSTatUmb3O7K4SspGPaEOQg=="],
"@rivet-dev/agentos/rivetkit/uuid": ["uuid@12.0.1", "", { "bin": { "uuid": "dist/bin/uuid" } }, "sha512-9obBF8sMIHJWNQaO6IGOG8giGa/jUpKX34bz6o4whVs8M0WAvhID2tNxYp6A2XEBJPuZSX8wsS/6TEKfIDc+nw=="],
"@rivetkit/framework-base/rivetkit/@rivetkit/engine-cli": ["@rivetkit/engine-cli@2.3.7", "", { "optionalDependencies": { "@rivetkit/engine-cli-darwin-arm64": "2.3.7", "@rivetkit/engine-cli-darwin-x64": "2.3.7", "@rivetkit/engine-cli-linux-arm64-musl": "2.3.7", "@rivetkit/engine-cli-linux-x64-musl": "2.3.7", "@rivetkit/engine-cli-win32-x64": "2.3.7" } }, "sha512-CezLwJ0B7dWDbA7qM6Aq04mwnrJAdrDActRrrcb4NBa20h7wO9KPAYBwyJz/dRnKm9EUUNcFZ6hrSDpp6T3+Rg=="],
@@ -4394,8 +4388,6 @@
"@rivetkit/framework-base/rivetkit/@rivetkit/workflow-engine": ["@rivetkit/workflow-engine@2.3.7", "", { "dependencies": { "@rivetkit/bare-ts": "^0.6.2", "cbor-x": "^1.6.0", "fdb-tuple": "^1.0.0", "pino": "^9.6.0", "vbare": "^0.0.4" } }, "sha512-+C4kGuSNysrw2zs/LQoYxouOxufKHHMxiIXKF8Ab9GU2BCJ32RUST9dYU/gaQE/w9uM+hqAHUz0DTlOBYU6p6g=="],
"@rivetkit/framework-base/rivetkit/hono": ["hono@4.12.31", "", {}, "sha512-zJIHFrl6bq3RDd2YusFNCDlM8qUprxKswyi/OPzPyzKDdyBXDqWx8bZlZ7R+saTdSTatUmb3O7K4SspGPaEOQg=="],
"@rivetkit/framework-base/rivetkit/uuid": ["uuid@12.0.1", "", { "bin": { "uuid": "dist/bin/uuid" } }, "sha512-9obBF8sMIHJWNQaO6IGOG8giGa/jUpKX34bz6o4whVs8M0WAvhID2tNxYp6A2XEBJPuZSX8wsS/6TEKfIDc+nw=="],
"@rivetkit/react/@tanstack/react-store/@tanstack/store": ["@tanstack/store@0.7.7", "", {}, "sha512-xa6pTan1bcaqYDS9BDpSiS63qa6EoDkPN9RsRaxHuDdVDNntzq3xNwR5YKTU/V3SkSyC9T4YVOPh2zRQN0nhIQ=="],
@@ -4414,8 +4406,6 @@
"@rivetkit/react/rivetkit/@rivetkit/workflow-engine": ["@rivetkit/workflow-engine@2.3.7", "", { "dependencies": { "@rivetkit/bare-ts": "^0.6.2", "cbor-x": "^1.6.0", "fdb-tuple": "^1.0.0", "pino": "^9.6.0", "vbare": "^0.0.4" } }, "sha512-+C4kGuSNysrw2zs/LQoYxouOxufKHHMxiIXKF8Ab9GU2BCJ32RUST9dYU/gaQE/w9uM+hqAHUz0DTlOBYU6p6g=="],
"@rivetkit/react/rivetkit/hono": ["hono@4.12.31", "", {}, "sha512-zJIHFrl6bq3RDd2YusFNCDlM8qUprxKswyi/OPzPyzKDdyBXDqWx8bZlZ7R+saTdSTatUmb3O7K4SspGPaEOQg=="],
"@rivetkit/react/rivetkit/uuid": ["uuid@12.0.1", "", { "bin": { "uuid": "dist/bin/uuid" } }, "sha512-9obBF8sMIHJWNQaO6IGOG8giGa/jUpKX34bz6o4whVs8M0WAvhID2tNxYp6A2XEBJPuZSX8wsS/6TEKfIDc+nw=="],
"@secure-exec/nodejs/esbuild/@esbuild/aix-ppc64": ["@esbuild/aix-ppc64@0.27.7", "", { "os": "aix", "cpu": "ppc64" }, "sha512-EKX3Qwmhz1eMdEJokhALr0YiD0lhQNwDqkPYyPhiSwKrh7/4KRjQc04sZ8db+5DVVnZ1LmbNDI1uAMPEUBnQPg=="],

View File

@@ -14,10 +14,12 @@
"dependencies": {
"@code/backend": "workspace:*",
"@code/env": "workspace:*",
"@code/primitives": "workspace:*",
"@flue/runtime": "latest",
"@rivet-dev/agentos-core": "catalog:",
"convex": "catalog:",
"hono": "4.12.30",
"effect": "catalog:",
"hono": "4.12.31",
"valibot": "^1.4.2"
},
"devDependencies": {

View File

@@ -3,6 +3,8 @@ import { registerProvider } from "@flue/runtime";
import { flue } from "@flue/runtime/routing";
import { Hono } from "hono";
import { projectRequestRoute } from "./project-request";
const agentEnv = parseAgentEnv(process.env);
registerProvider(agentEnv.AGENT_MODEL_PROVIDER, {
@@ -18,6 +20,7 @@ registerProvider(agentEnv.AGENT_MODEL_PROVIDER, {
});
const app = new Hono();
app.post("/project-requests", projectRequestRoute);
app.route("/", flue());
export default app;

View File

@@ -98,7 +98,9 @@ const { CONVEX_URL } = parseAgentEnv(process.env);
* caller's JWT is set on this instance only and discarded when the request
* ends. We never mutate auth on a shared/global Convex client.
*/
const createAuthenticatedClient = (accessToken: string): ConvexHttpClient => {
export const createAuthenticatedClient = (
accessToken: string
): ConvexHttpClient => {
const client = new ConvexHttpClient(CONVEX_URL);
client.setAuth(accessToken);
return client;
@@ -120,7 +122,7 @@ interface HeaderSource {
* Extract and validate the Bearer access token from the request. Returns null
* when absent or malformed so the caller can produce a clean 401.
*/
const extractBearerToken = (request: HeaderSource): string | null => {
export const extractBearerToken = (request: HeaderSource): string | null => {
const header = request.headers.get("authorization");
if (!header) {
return null;

View File

@@ -0,0 +1,162 @@
import {
decodeProjectIssueRequest,
ProjectIssueDispatchInput,
ProjectIssueRequestError,
ProjectIssueRequestResult,
ProjectIssueValidationError,
} from "@code/primitives/project-issue";
import type { ProjectIssueRequest } from "@code/primitives/project-issue";
import { dispatch } from "@flue/runtime";
import type { ConvexHttpClient } from "convex/browser";
import { makeFunctionReference } from "convex/server";
import { Effect, Schema } from "effect";
import type { Context } from "hono";
import { createAuthenticatedClient, extractBearerToken } from "./auth";
interface ProjectIssueCreateArgs extends Record<string, unknown> {
readonly body: string;
readonly projectId: string;
readonly title: string;
}
const createIssue = makeFunctionReference<
"mutation",
ProjectIssueCreateArgs,
string
>("projectIssues:create");
const createIssueFromSignal = makeFunctionReference<
"mutation",
{ readonly signalId: string },
{ readonly issueId: string; readonly projectId: string }
>("projectIssues:createFromSignal");
const beginIssue = makeFunctionReference<
"mutation",
{ readonly issueId: string },
"queued" | "working"
>("projectIssues:begin");
const markDispatchFailed = makeFunctionReference<
"mutation",
{ readonly error: string; readonly issueId: string },
null
>("projectIssues:markDispatchFailed");
const invalidRequest = (message: string) =>
Response.json({ error: message }, { status: 400 });
const knownAuthorizationFailure = (message: string): boolean =>
/authentication required|membership required|project not found|signal not found|not project-scoped/iu.test(
message
);
const errorMessage = (error: unknown): string =>
error instanceof Error ? error.message : String(error);
const decodeRequest = async (input: unknown): Promise<ProjectIssueRequest> => {
try {
return await Effect.runPromise(decodeProjectIssueRequest(input));
} catch (error) {
if (
error instanceof ProjectIssueRequestError ||
error instanceof ProjectIssueValidationError
) {
throw invalidRequest(
error instanceof Error ? error.message : "Invalid project request"
);
}
throw error;
}
};
const createIssueForRequest = async (
client: ConvexHttpClient,
request: Awaited<ReturnType<typeof decodeRequest>>
): Promise<{ readonly issueId: string; readonly projectId: string }> => {
if (request.kind === "signal") {
return client.mutation(createIssueFromSignal, {
signalId: request.signalId,
});
}
const issueId = await client.mutation(createIssue, {
body: request.body,
projectId: request.projectId,
title: request.title,
});
return { issueId, projectId: request.projectId };
};
export const projectRequestRoute = async (c: Context): Promise<Response> => {
const accessToken = extractBearerToken(c.req.raw);
if (!accessToken) {
return c.json({ error: "Unauthorized" }, 401);
}
let input: unknown;
try {
input = await c.req.json();
} catch {
return invalidRequest("Request body must be valid JSON");
}
let request: Awaited<ReturnType<typeof decodeRequest>>;
try {
request = await decodeRequest(input);
} catch (error) {
if (error instanceof Response) {
return error;
}
return c.json({ error: "Invalid project request" }, 400);
}
const client = createAuthenticatedClient(accessToken);
let issue:
| { readonly issueId: string; readonly projectId: string }
| undefined;
try {
issue = await createIssueForRequest(client, request);
const status = await client.mutation(beginIssue, {
issueId: issue.issueId,
});
const dispatchInput = Schema.decodeUnknownSync(ProjectIssueDispatchInput)({
issueId: issue.issueId,
kind: "project.issue.started",
projectId: issue.projectId,
});
const receipt = await dispatch({
agent: "project-manager",
id: issue.issueId,
input: dispatchInput,
});
const result = Schema.decodeUnknownSync(ProjectIssueRequestResult)({
acceptedAt: receipt.acceptedAt,
dispatchId: receipt.dispatchId,
issueId: issue.issueId,
projectId: issue.projectId,
status,
});
return c.json(result, 202);
} catch (error) {
const message = errorMessage(error);
if (issue) {
try {
await client.mutation(markDispatchFailed, {
error: message,
issueId: issue.issueId,
});
} catch {
// Preserve the original request failure; the issue remains inspectable.
}
}
return c.json(
{
error: knownAuthorizationFailure(message)
? "Project request is not authorized"
: "Project request could not be dispatched",
},
knownAuthorizationFailure(message) ? 403 : 502
);
}
};

View File

@@ -0,0 +1,177 @@
import { type TestConvex, convexTest } from "convex-test";
import { anyApi } from "convex/server";
import { describe, expect, test } from "vitest";
import { internal } from "./_generated/api";
import type { Id } from "./_generated/dataModel";
import schema from "./schema";
declare global {
interface ImportMeta {
readonly glob: (pattern: string) => Record<string, () => Promise<unknown>>;
}
}
const modules = import.meta.glob("./**/*.ts");
const api = anyApi;
const ID_A = "https://convex.test|project-request-a";
const ID_B = "https://convex.test|project-request-b";
const identityA = { tokenIdentifier: ID_A };
const identityB = { tokenIdentifier: ID_B };
const createProject = async (
t: TestConvex<typeof schema>
): Promise<Id<"projects">> => {
await t
.withIdentity(identityA)
.mutation(api.organizations.ensurePersonalOrganization, {});
const outcome = await t
.withIdentity(identityA)
.mutation(internal.projects.persistPublicGitImport, {
userId: ID_A,
source: {
host: "github.com",
normalizedUrl: "https://github.com/test/project-request",
projectName: "project-request",
repositoryPath: "test/project-request",
url: "https://github.com/test/project-request",
},
remote: {
defaultBranch: "main",
documents: [
{
content: "# Project Request\n",
kind: "readme" as const,
path: "README.md",
},
],
warnings: [],
},
});
return outcome.id as unknown as Id<"projects">;
};
describe("project issue request smoke", () => {
test("project Signal becomes a queued project issue with durable evidence", async () => {
const t = convexTest({ schema, modules });
const projectId = await createProject(t);
const organization = await t
.withIdentity(identityA)
.mutation(api.organizations.ensurePersonalOrganization, {});
const begun = await t
.withIdentity(identityA)
.mutation(api.conversationMessages.beginUserMessage, {
clientRequestId: "project-request-1",
organizationId: organization._id,
rawText: "Staging deploys fail during DNS resolution.",
});
await t
.withIdentity(identityA)
.mutation(api.conversationMessages.markAdmitted, {
clientRequestId: "project-request-1",
organizationId: organization._id,
submissionId: "submission-project-request-1",
});
const signal = await t
.withIdentity(identityA)
.mutation(api.signals.createFromMessages, {
conversationId: organization._id,
messageIds: [begun.messageId],
organizationId: organization._id,
problemStatement: {
constraints: ["Do not change production"],
desiredOutcome: "Staging deploys successfully",
summary: "The staging deploy fails during DNS resolution.",
title: "Fix the staging deploy DNS failure",
},
processedBy: {
agentInstanceId: String(organization._id),
agentName: "zopu",
},
projectId,
});
const created = await t
.withIdentity(identityA)
.mutation(api.projectIssues.createFromSignal, {
signalId: signal.signalId,
});
const issue = await t.run(async (ctx) => ctx.db.get(created.issueId));
expect(issue).toMatchObject({
body: expect.stringContaining("Source Signal"),
projectId,
status: "open",
title: "Fix the staging deploy DNS failure",
});
const status = await t
.withIdentity(identityA)
.mutation(api.projectIssues.begin, { issueId: created.issueId });
expect(status).toBe("queued");
const events = await t.run(async (ctx) =>
ctx.db
.query("projectEvents")
.withIndex("by_project_and_createdAt", (q) =>
q.eq("projectId", projectId)
)
.collect()
);
expect(events.map((event) => event.kind)).toEqual([
"project.connected",
"issue.created",
"issue.queued",
]);
});
test("a different organization cannot consume the project Signal", async () => {
const t = convexTest({ schema, modules });
const projectId = await createProject(t);
await t
.withIdentity(identityB)
.mutation(api.organizations.ensurePersonalOrganization, {});
const organization = await t
.withIdentity(identityA)
.mutation(api.organizations.ensurePersonalOrganization, {});
const begun = await t
.withIdentity(identityA)
.mutation(api.conversationMessages.beginUserMessage, {
clientRequestId: "project-request-cross-org",
organizationId: organization._id,
rawText: "Private project request",
});
await t
.withIdentity(identityA)
.mutation(api.conversationMessages.markAdmitted, {
clientRequestId: "project-request-cross-org",
organizationId: organization._id,
submissionId: "submission-project-request-cross-org",
});
const signal = await t
.withIdentity(identityA)
.mutation(api.signals.createFromMessages, {
conversationId: organization._id,
messageIds: [begun.messageId],
organizationId: organization._id,
problemStatement: {
constraints: [],
desiredOutcome: "Keep the request private",
summary: "This request belongs to the first organization.",
title: "Private project request",
},
processedBy: {
agentInstanceId: String(organization._id),
agentName: "zopu",
},
projectId,
});
await expect(
t.withIdentity(identityB).mutation(api.projectIssues.createFromSignal, {
signalId: signal.signalId,
})
).rejects.toThrow(/membership required/u);
});
});

View File

@@ -1,8 +1,82 @@
import {
projectIssueDraftFromSignal,
queueProjectIssue,
validateProjectIssueDraft,
type ProjectIssueDraft,
} from "@code/primitives/project-issue";
import { ConvexError, v } from "convex/values";
import { Effect } from "effect";
import { mutation, query } from "./_generated/server";
import type { Id } from "./_generated/dataModel";
import { mutation, type MutationCtx, query } from "./_generated/server";
import { requireProjectMember } from "./authz";
const toDomainError = (error: unknown) => {
return new ConvexError(
error instanceof Error ? error.message : "Invalid project issue"
);
};
const decodeDraft = async (input: unknown): Promise<ProjectIssueDraft> => {
try {
return await Effect.runPromise(validateProjectIssueDraft(input));
} catch (error) {
throw toDomainError(error);
}
};
const draftFromSignal = async (input: unknown): Promise<ProjectIssueDraft> => {
try {
return await Effect.runPromise(projectIssueDraftFromSignal(input));
} catch (error) {
throw toDomainError(error);
}
};
const insertIssue = async (
ctx: MutationCtx,
projectId: Id<"projects">,
draft: ProjectIssueDraft
): Promise<Id<"projectIssues">> => {
const latest = await ctx.db
.query("projectIssues")
.withIndex("by_project_and_number", (q) => q.eq("projectId", projectId))
.order("desc")
.first();
const timestamp = Date.now();
const number = (latest?.number ?? 0) + 1;
const issueId = await ctx.db.insert("projectIssues", {
body: draft.body,
createdAt: timestamp,
number,
projectId,
status: "open",
title: draft.title,
updatedAt: timestamp,
});
const workArtifact = await ctx.db
.query("projectArtifacts")
.withIndex("by_project_and_path", (q) =>
q.eq("projectId", projectId).eq("path", "work.md")
)
.unique();
if (workArtifact) {
await ctx.db.patch("projectArtifacts", workArtifact._id, {
content: `${workArtifact.content}\n## Issue ${number}: ${draft.title}\n\nStatus: open\n\n${draft.body}\n`,
revision: workArtifact.revision + 1,
updatedAt: timestamp,
});
}
await ctx.db.insert("projectEvents", {
createdAt: timestamp,
data: { number, title: draft.title },
issueId,
kind: "issue.created",
projectId,
});
return issueId;
};
export const create = mutation({
args: {
body: v.string(),
@@ -11,55 +85,36 @@ export const create = mutation({
},
handler: async (ctx, args) => {
await requireProjectMember(ctx, args.projectId);
const title = args.title.trim();
const body = args.body.trim();
if (title.length < 3 || title.length > 160) {
throw new ConvexError("Issue title must be between 3 and 160 characters");
const draft = await decodeDraft({ body: args.body, title: args.title });
return insertIssue(ctx, args.projectId, draft);
},
});
export const createFromSignal = mutation({
args: { signalId: v.id("signals") },
handler: async (ctx, args) => {
const signal = await ctx.db.get("signals", args.signalId);
if (!signal) {
throw new ConvexError("Signal not found");
}
if (body.length < 10 || body.length > 10_000) {
if (!signal.projectId) {
throw new ConvexError("Signal is not project-scoped");
}
const { organizationId } = await requireProjectMember(
ctx,
signal.projectId
);
if (signal.organizationId !== organizationId) {
throw new ConvexError(
"Issue description must be between 10 and 10000 characters"
"Signal does not belong to the project organization"
);
}
const latest = await ctx.db
.query("projectIssues")
.withIndex("by_project_and_number", (q) =>
q.eq("projectId", args.projectId)
)
.order("desc")
.first();
const timestamp = Date.now();
const number = (latest?.number ?? 0) + 1;
const issueId = await ctx.db.insert("projectIssues", {
body,
createdAt: timestamp,
number,
projectId: args.projectId,
status: "open",
title,
updatedAt: timestamp,
const draft = await draftFromSignal({
problemStatement: signal.problemStatement,
signalId: String(signal._id),
});
const workArtifact = await ctx.db
.query("projectArtifacts")
.withIndex("by_project_and_path", (q) =>
q.eq("projectId", args.projectId).eq("path", "work.md")
)
.unique();
if (workArtifact) {
await ctx.db.patch("projectArtifacts", workArtifact._id, {
content: `${workArtifact.content}\n## Issue ${number}: ${title}\n\nStatus: open\n\n${body}\n`,
revision: workArtifact.revision + 1,
updatedAt: timestamp,
});
}
await ctx.db.insert("projectEvents", {
createdAt: timestamp,
data: { number, title },
issueId,
kind: "issue.created",
projectId: args.projectId,
});
return issueId;
const issueId = await insertIssue(ctx, signal.projectId, draft);
return { issueId, projectId: signal.projectId };
},
});
@@ -71,12 +126,18 @@ export const begin = mutation({
throw new ConvexError("Issue not found");
}
await requireProjectMember(ctx, issue.projectId);
if (issue.status === "queued" || issue.status === "working") {
return issue.status;
let status: "queued" | "working";
try {
status = await Effect.runPromise(queueProjectIssue(issue.status));
} catch (error) {
throw toDomainError(error);
}
if (status === issue.status) {
return status;
}
const timestamp = Date.now();
await ctx.db.patch("projectIssues", issue._id, {
status: "queued",
status,
updatedAt: timestamp,
});
await ctx.db.insert("projectEvents", {

View File

@@ -7,6 +7,7 @@
".": "./src/index.ts",
"./agent-os": "./src/agent-os.ts",
"./project": "./src/project.ts",
"./project-issue": "./src/project-issue.ts",
"./signal": "./src/signal.ts"
},
"scripts": {

View File

@@ -1,4 +1,5 @@
// oxlint-disable-next-line no-barrel-file -- The package root intentionally exposes its public modules.
export * from "./agent-os";
export * from "./project";
export * from "./project-issue";
export * from "./signal";

View File

@@ -0,0 +1,77 @@
import { Effect } from "effect";
import { describe, expect, it } from "vitest";
import {
decodeProjectIssueRequest,
ProjectIssueTransitionError,
ProjectIssueValidationError,
projectIssueDraftFromSignal,
queueProjectIssue,
validateProjectIssueDraft,
} from "./project-issue";
describe("project issue primitives", () => {
it("normalizes an explicit project request", () => {
const request = Effect.runSync(
decodeProjectIssueRequest({
body: " Investigate the staging failure. ",
kind: "explicit",
projectId: "project-1",
title: " Staging failure ",
})
);
expect(request).toEqual({
body: "Investigate the staging failure.",
kind: "explicit",
projectId: "project-1",
title: "Staging failure",
});
});
it("rejects invalid issue drafts before persistence", () => {
const error = Effect.runSync(
Effect.flip(
validateProjectIssueDraft({
body: "too short",
title: "No",
})
)
);
expect(error).toBeInstanceOf(ProjectIssueValidationError);
expect(error.reason).toBe("TitleTooShort");
});
it("projects a project Signal into the existing issue shape", () => {
const draft = Effect.runSync(
projectIssueDraftFromSignal({
problemStatement: {
constraints: ["Do not change production"],
desiredOutcome: "Staging deploys successfully",
summary: "The deploy fails during DNS resolution",
title: "Fix staging deploy DNS failure",
},
signalId: "signal-1",
})
);
expect(draft.title).toBe("Fix staging deploy DNS failure");
expect(draft.body).toContain("The deploy fails during DNS resolution");
expect(draft.body).toContain("Source Signal: signal-1");
});
it("queues retryable issue states and preserves active states", () => {
expect(Effect.runSync(queueProjectIssue("open"))).toBe("queued");
expect(Effect.runSync(queueProjectIssue("failed"))).toBe("queued");
expect(Effect.runSync(queueProjectIssue("needs-input"))).toBe("queued");
expect(Effect.runSync(queueProjectIssue("working"))).toBe("working");
});
it("does not reopen completed issues", () => {
const error = Effect.runSync(Effect.flip(queueProjectIssue("completed")));
expect(error).toBeInstanceOf(ProjectIssueTransitionError);
expect(error.reason).toBe("CompletedIssue");
});
});

View File

@@ -0,0 +1,265 @@
/* eslint-disable max-classes-per-file -- each domain failure has a distinct tagged reason. */
import { Effect, Schema } from "effect";
import { ProblemStatement, ProjectId, SignalId } from "./signal.js";
const MeaningfulString = Schema.String.check(
Schema.makeFilter((value) => value.trim().length > 0, {
expected: "a non-empty string",
})
);
export const ProjectIssueId = MeaningfulString.pipe(
Schema.brand("ProjectIssueId")
);
export type ProjectIssueId = typeof ProjectIssueId.Type;
export const ProjectIssueStatus = Schema.Literals([
"open",
"queued",
"working",
"needs-input",
"completed",
"failed",
]);
export type ProjectIssueStatus = typeof ProjectIssueStatus.Type;
const ProjectIssueDraftInput = Schema.Struct({
body: Schema.String,
title: Schema.String,
});
export const ProjectIssueDraft = Schema.Struct({
body: MeaningfulString,
title: MeaningfulString,
});
export type ProjectIssueDraft = typeof ProjectIssueDraft.Type;
export const ExplicitProjectIssueRequest = Schema.Struct({
body: Schema.String,
kind: Schema.Literal("explicit"),
projectId: ProjectId,
title: Schema.String,
});
export const SignalProjectIssueRequest = Schema.Struct({
kind: Schema.Literal("signal"),
signalId: SignalId,
});
export const ProjectIssueRequest = Schema.Union([
ExplicitProjectIssueRequest,
SignalProjectIssueRequest,
]);
export type ProjectIssueRequest = typeof ProjectIssueRequest.Type;
export const ProjectIssueDispatchInput = Schema.Struct({
issueId: ProjectIssueId,
kind: Schema.Literal("project.issue.started"),
projectId: ProjectId,
});
export type ProjectIssueDispatchInput = typeof ProjectIssueDispatchInput.Type;
export const ProjectIssueRequestResult = Schema.Struct({
acceptedAt: MeaningfulString,
dispatchId: MeaningfulString,
issueId: ProjectIssueId,
projectId: ProjectId,
status: Schema.Literals(["queued", "working"]),
});
export type ProjectIssueRequestResult = typeof ProjectIssueRequestResult.Type;
export const ProjectIssueValidationReason = Schema.Literals([
"InvalidInput",
"TitleTooShort",
"TitleTooLong",
"BodyTooShort",
"BodyTooLong",
]);
export type ProjectIssueValidationReason =
typeof ProjectIssueValidationReason.Type;
export class ProjectIssueValidationError extends Schema.TaggedErrorClass<ProjectIssueValidationError>()(
"ProjectIssueValidationError",
{
message: Schema.String,
reason: ProjectIssueValidationReason,
}
) {}
export const ProjectIssueRequestErrorReason = Schema.Literals([
"InvalidInput",
"InvalidProjectId",
"InvalidSignalId",
]);
export type ProjectIssueRequestErrorReason =
typeof ProjectIssueRequestErrorReason.Type;
export class ProjectIssueRequestError extends Schema.TaggedErrorClass<ProjectIssueRequestError>()(
"ProjectIssueRequestError",
{
message: Schema.String,
reason: ProjectIssueRequestErrorReason,
}
) {}
export const ProjectIssueTransitionReason = Schema.Literals([
"InvalidStatus",
"CompletedIssue",
]);
export type ProjectIssueTransitionReason =
typeof ProjectIssueTransitionReason.Type;
export class ProjectIssueTransitionError extends Schema.TaggedErrorClass<ProjectIssueTransitionError>()(
"ProjectIssueTransitionError",
{
message: Schema.String,
reason: ProjectIssueTransitionReason,
}
) {}
export const validateProjectIssueDraft = (
input: unknown
): Effect.Effect<ProjectIssueDraft, ProjectIssueValidationError> =>
Effect.gen(function* validateDraft() {
const decoded = yield* Schema.decodeUnknownEffect(ProjectIssueDraftInput)(
input
).pipe(
Effect.mapError(
() =>
new ProjectIssueValidationError({
message: "Issue title and body are required",
reason: "InvalidInput",
})
)
);
const title = decoded.title.trim();
const body = decoded.body.trim();
if (title.length < 3) {
return yield* Effect.fail(
new ProjectIssueValidationError({
message: "Issue title must be between 3 and 160 characters",
reason: "TitleTooShort",
})
);
}
if (title.length > 160) {
return yield* Effect.fail(
new ProjectIssueValidationError({
message: "Issue title must be between 3 and 160 characters",
reason: "TitleTooLong",
})
);
}
if (body.length < 10) {
return yield* Effect.fail(
new ProjectIssueValidationError({
message: "Issue description must be between 10 and 10000 characters",
reason: "BodyTooShort",
})
);
}
if (body.length > 10_000) {
return yield* Effect.fail(
new ProjectIssueValidationError({
message: "Issue description must be between 10 and 10000 characters",
reason: "BodyTooLong",
})
);
}
return { body, title };
});
export const decodeProjectIssueRequest = (
input: unknown
): Effect.Effect<
ProjectIssueRequest,
ProjectIssueRequestError | ProjectIssueValidationError
> =>
Effect.gen(function* decodeRequest() {
const request = yield* Schema.decodeUnknownEffect(ProjectIssueRequest)(
input
).pipe(
Effect.mapError(
() =>
new ProjectIssueRequestError({
message:
"Use a signal request or an explicit request with projectId, title, and body",
reason: "InvalidInput",
})
)
);
if (request.kind === "explicit") {
const draft = yield* validateProjectIssueDraft(request);
return { ...request, ...draft };
}
return request;
});
export const projectIssueDraftFromSignal = (input: unknown) =>
Effect.gen(function* draftFromSignal() {
const source = yield* Schema.decodeUnknownEffect(
Schema.Struct({
problemStatement: ProblemStatement,
signalId: SignalId,
})
)(input).pipe(
Effect.mapError(
() =>
new ProjectIssueRequestError({
message: "Project Signal is missing a valid problem statement",
reason: "InvalidSignalId",
})
)
);
const constraints =
source.problemStatement.constraints.length === 0
? "None"
: source.problemStatement.constraints
.map((constraint) => `- ${constraint}`)
.join("\n");
return yield* validateProjectIssueDraft({
body: [
source.problemStatement.summary,
`Desired outcome: ${source.problemStatement.desiredOutcome}`,
`Constraints:\n${constraints}`,
`Source Signal: ${source.signalId}`,
].join("\n\n"),
title: source.problemStatement.title,
});
});
export const queueProjectIssue = (
input: unknown
): Effect.Effect<"queued" | "working", ProjectIssueTransitionError> =>
Effect.gen(function* queueIssue() {
const status = yield* Schema.decodeUnknownEffect(ProjectIssueStatus)(
input
).pipe(
Effect.mapError(
() =>
new ProjectIssueTransitionError({
message: "Issue has an invalid status",
reason: "InvalidStatus",
})
)
);
if (status === "queued" || status === "working") {
return status;
}
if (status === "completed") {
return yield* Effect.fail(
new ProjectIssueTransitionError({
message: "Completed issues cannot be queued again",
reason: "CompletedIssue",
})
);
}
return "queued" as const;
});