Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
-- AlterTable
ALTER TABLE "tasks" ADD COLUMN "audit_callback_key" TEXT;
4 changes: 4 additions & 0 deletions prisma/schema.prisma
Original file line number Diff line number Diff line change
Expand Up @@ -699,6 +699,10 @@ model Task {
// and is verified against this secret on receipt (mirrors
// `agentWebhookSecret`, kept separate so the names stay honest).
codeChangeWebhookSecret String? @map("code_change_webhook_secret")
// Encrypted one-time key authenticating the pod-resident auditor's verdict
// callback (mirrors `agentPassword`/`codeChangeWebhookSecret`). Cleared once
// the verdict lands.
auditCallbackKey String? @map("audit_callback_key")
agentLogs AgentLog[]
chatMessages ChatMessage[]
deployments Deployment[]
Expand Down
39 changes: 39 additions & 0 deletions src/app/api/ask/quick/route.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import { NextRequest, NextResponse, after } from "next/server";
import { randomUUID } from "crypto";
import { validationError, serverError, forbiddenError, isApiError } from "@/types/errors";
import { validateUserBelongsToOrg, validateWorkspaceAccess } from "@/services/workspace";
import { ModelMessage, createUIMessageStream, createUIMessageStreamResponse, toUIMessageStream } from "ai";
Expand Down Expand Up @@ -650,6 +651,44 @@ export async function POST(request: NextRequest) {
);
writerRef.write({ type: "data-usage", data: cumulativeUsage });
}

// Detect a verification verdict (submit_verdict from the verify
// MCP tools) and emit it as a chat artifact for the display.
for (const part of (sf.content ?? []) as Array<Record<string, unknown>>) {
if (
part?.type !== "tool-result" ||
typeof part.toolName !== "string" ||
!(part.toolName === "submit_verdict" || part.toolName.endsWith("_submit_verdict"))
) {
continue;
}
let verdict: unknown = part.output ?? part.result;
if (typeof verdict === "string") {
try {
verdict = JSON.parse(verdict);
} catch {
/* leave as-is */
}
}
const wrapped = verdict as { content?: Array<{ type?: string; text?: string }> };
if (Array.isArray(wrapped?.content)) {
const t = wrapped.content.find((c) => c?.type === "text")?.text;
if (t) {
try {
verdict = JSON.parse(t);
} catch {
/* keep */
}
}
}
const v = verdict as { overall?: unknown; claims?: unknown };
if (v && typeof v === "object" && "overall" in v && "claims" in v) {
writerRef.write({
type: "data-artifact",
data: { id: randomUUID(), type: "VERIFY", data: verdict },
});
}
}
},
onFinish: async ({
usage,
Expand Down
142 changes: 142 additions & 0 deletions src/app/api/tasks/[taskId]/audit/callback/route.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,142 @@
import { NextRequest, NextResponse } from "next/server";
import { db } from "@/lib/db";
import { EncryptionService, timingSafeEqual } from "@/lib/encryption";
import { ChatRole, ChatStatus, ArtifactType, WorkflowStatus, NotificationTriggerType, Prisma } from "@prisma/client";
import { createAndSendNotification } from "@/services/notifications";
import { pusherServer, getTaskChannelName, PUSHER_EVENTS } from "@/lib/pusher";
import type { AuditVerdict } from "@/services/auditor/types";

export const fetchCache = "force-no-store";

const encryptionService = EncryptionService.getInstance();

export async function POST(request: NextRequest, { params }: { params: Promise<{ taskId: string }> }) {
try {
const { taskId } = await params;
if (!taskId) {
return NextResponse.json({ error: "Task ID required" }, { status: 400 });
}

const apiKey = request.headers.get("x-api-key");
if (!apiKey) {
return NextResponse.json({ error: "Unauthorized" }, { status: 401 });
}

const task = await db.task.findUnique({
where: { id: taskId, deleted: false },
select: {
id: true,
auditCallbackKey: true,
workspaceId: true,
assigneeId: true,
createdById: true,
title: true,
workspace: { select: { slug: true } },
},
});

if (!task) {
return NextResponse.json({ error: "Task not found" }, { status: 404 });
}

if (!task.auditCallbackKey) {
return NextResponse.json({ error: "Unauthorized" }, { status: 401 });
}

let decryptedKey: string;
try {
const encryptedData = JSON.parse(task.auditCallbackKey);
decryptedKey = encryptionService.decryptField("auditCallbackKey", encryptedData);
} catch (error) {
console.error("Failed to decrypt auditCallbackKey:", error);
return NextResponse.json({ error: "Unauthorized" }, { status: 401 });
}

if (!timingSafeEqual(apiKey, decryptedKey)) {
return NextResponse.json({ error: "Unauthorized" }, { status: 401 });
}

let verdict: AuditVerdict;
try {
verdict = (await request.json()) as AuditVerdict;
} catch {
return NextResponse.json({ error: "Invalid JSON body" }, { status: 400 });
}

const overall = verdict.overall;
const works = overall === "works";
const workflowStatus = works ? WorkflowStatus.COMPLETED : WorkflowStatus.FAILED;

const summaryText = verdict.summary || `Audit verdict: ${overall}`;

const evidence = Array.isArray(verdict.evidence) ? verdict.evidence : [];

await db.chatMessage.create({
data: {
taskId,
message: summaryText,
role: ChatRole.ASSISTANT,
status: ChatStatus.SENT,
artifacts: {
create: [
{
type: ArtifactType.VERIFY,
content: {
overall,
claims: verdict.claims ?? [],
observations: verdict.observations ?? [],
summary: summaryText,
evidence,
startedAt: verdict.startedAt ?? "",
finishedAt: verdict.finishedAt ?? "",
error: verdict.error ?? "",
} as unknown as Prisma.InputJsonValue,
icon: "verify",
},
],
},
},
});

await db.task.update({
where: { id: taskId },
data: {
workflowStatus,
workflowCompletedAt: new Date(),
auditCallbackKey: null,
},
});

if (!works) {
const targetUserId = task.assigneeId ?? task.createdById;
const baseUrl = process.env.NEXTAUTH_URL || "http://localhost:3000";
const taskUrl = `${baseUrl}/w/${task.workspace.slug}/task/${taskId}`;
try {
await createAndSendNotification({
targetUserId,
taskId,
workspaceId: task.workspaceId,
notificationType: NotificationTriggerType.WORKFLOW_HALTED,
message: `Audit flagged task '${task.title}' as ${overall} — needs a human: ${taskUrl}`,
});
} catch (notifError) {
console.error("[auditor:callback] Failed to fire WORKFLOW_HALTED notification:", notifError);
}
}

try {
await pusherServer.trigger(getTaskChannelName(taskId), PUSHER_EVENTS.WORKFLOW_STATUS_UPDATE, {
taskId,
workflowStatus,
timestamp: new Date(),
});
} catch (pusherError) {
console.error("[auditor:callback] Pusher broadcast failed:", pusherError);
}

return NextResponse.json({ success: true }, { status: 200 });
} catch (error) {
console.error("Unexpected error in audit callback:", error);
return NextResponse.json({ error: "Internal error" }, { status: 500 });
}
}
45 changes: 45 additions & 0 deletions src/app/api/tasks/[taskId]/audit/route.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,45 @@
import { NextRequest, NextResponse } from "next/server";
import { getMiddlewareContext, requireAuth } from "@/lib/middleware/utils";
import { db } from "@/lib/db";
import { validateWorkspaceAccessById } from "@/services/workspace";
import { startAudit } from "@/services/auditor/trigger";

export const fetchCache = "force-no-store";

export async function POST(request: NextRequest, { params }: { params: Promise<{ taskId: string }> }) {
try {
const context = getMiddlewareContext(request);
const userOrResponse = requireAuth(context);
if (userOrResponse instanceof NextResponse) return userOrResponse;

const { taskId } = await params;
if (!taskId) {
return NextResponse.json({ error: "Task ID required" }, { status: 400 });
}

const task = await db.task.findUnique({
where: { id: taskId, deleted: false },
select: { id: true, workspaceId: true },
});

if (!task) {
return NextResponse.json({ error: "Task not found" }, { status: 404 });
}

const access = await validateWorkspaceAccessById(task.workspaceId, userOrResponse.id);
if (!access.hasAccess) {
return NextResponse.json({ error: "Forbidden" }, { status: 403 });
}

const result = await startAudit(taskId, userOrResponse.id);

if (!result.success) {
return NextResponse.json({ error: result.error ?? "Failed to start audit" }, { status: 502 });
}

return NextResponse.json({ success: true, data: result }, { status: 202 });
} catch (error) {
console.error("Error starting audit:", error);
return NextResponse.json({ error: "Internal error" }, { status: 500 });
}
}
8 changes: 7 additions & 1 deletion src/app/org/[githubLogin]/_components/SidebarChat.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,8 @@ import {
Split,
X,
} from "lucide-react";
import { VerdictPill, isAuditVerdict } from "@/app/w/[slug]/task/[...taskParams]/artifacts/verdict";
import type { Artifact } from "@/lib/chat";
import { useSpeechRecognition } from "@/hooks/useSpeechRecognition";
import { useControlKeyHold } from "@/hooks/useControlKeyHold";
import { useVoiceCorrectionCapture } from "@/hooks/useVoiceCorrectionCapture";
Expand Down Expand Up @@ -727,8 +729,12 @@ function MessageArtifacts({ artifactIds }: { artifactIds?: string[] }) {
return (
<div className="space-y-1.5">
{artifacts.map((artifact) => {
if (artifact.type === "VERIFY" && isAuditVerdict(artifact.data)) {
return (
<VerdictPill key={artifact.id} artifact={{ content: artifact.data } as Artifact} />
);
}
// Unknown artifact type — render nothing rather than crash.
void artifact;
return null;
})}
</div>
Expand Down
20 changes: 20 additions & 0 deletions src/app/org/[githubLogin]/_state/useSendCanvasChatMessage.ts
Original file line number Diff line number Diff line change
Expand Up @@ -513,6 +513,26 @@ export function useSendCanvasChatMessage() {
}
}

// Register any mid-stream chat artifacts (e.g. verification verdicts)
// and attach their ids to the latest assistant message so
// MessageArtifacts renders them in the chat.
if (updatedMessage.artifacts?.length) {
const registerArtifact = useCanvasChatStore.getState().registerArtifact;
for (const art of updatedMessage.artifacts) {
registerArtifact({ id: art.id, type: art.type, conversationId, messageId, data: art.data });
for (let i = timelineMessages.length - 1; i >= 0; i--) {
const m = timelineMessages[i];
if (m.role === "assistant") {
const ids = m.artifactIds ?? [];
if (!ids.includes(art.id)) {
timelineMessages[i] = { ...m, artifactIds: [...ids, art.id] };
}
break;
}
}
}
}

const lastMsg = timelineMessages[timelineMessages.length - 1];
if (lastMsg?.toolCalls && lastMsg.toolCalls.length > 0) {
setActiveToolCalls(conversationId, lastMsg.toolCalls);
Expand Down
1 change: 1 addition & 0 deletions src/app/w/[slug]/task/[...taskParams]/artifacts/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,3 +10,4 @@ export { PublishScriptArtifact } from "./publish-script";
export { PublishPromptArtifact } from "./publish-prompt";
export { PublishSkillArtifact } from "./publish-skill";
export { BountyArtifact } from "./bounty";
export { VerdictArtifact, VerdictPill, isAuditVerdict } from "./verdict";
Loading
Loading