diff --git a/prisma/migrations/20260902000000_add_task_audit_callback_key/migration.sql b/prisma/migrations/20260902000000_add_task_audit_callback_key/migration.sql new file mode 100644 index 0000000000..abc85df6f3 --- /dev/null +++ b/prisma/migrations/20260902000000_add_task_audit_callback_key/migration.sql @@ -0,0 +1,2 @@ +-- AlterTable +ALTER TABLE "tasks" ADD COLUMN "audit_callback_key" TEXT; diff --git a/prisma/schema.prisma b/prisma/schema.prisma index bdb8e0e04b..5e27ccd13e 100644 --- a/prisma/schema.prisma +++ b/prisma/schema.prisma @@ -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[] diff --git a/src/app/api/ask/quick/route.ts b/src/app/api/ask/quick/route.ts index 0673cf40dc..d46474f2d4 100644 --- a/src/app/api/ask/quick/route.ts +++ b/src/app/api/ask/quick/route.ts @@ -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"; @@ -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>) { + 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, diff --git a/src/app/api/tasks/[taskId]/audit/callback/route.ts b/src/app/api/tasks/[taskId]/audit/callback/route.ts new file mode 100644 index 0000000000..49d8c756dc --- /dev/null +++ b/src/app/api/tasks/[taskId]/audit/callback/route.ts @@ -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 }); + } +} diff --git a/src/app/api/tasks/[taskId]/audit/route.ts b/src/app/api/tasks/[taskId]/audit/route.ts new file mode 100644 index 0000000000..a6b558c260 --- /dev/null +++ b/src/app/api/tasks/[taskId]/audit/route.ts @@ -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 }); + } +} diff --git a/src/app/org/[githubLogin]/_components/SidebarChat.tsx b/src/app/org/[githubLogin]/_components/SidebarChat.tsx index 58aae69c76..f90e487ca8 100644 --- a/src/app/org/[githubLogin]/_components/SidebarChat.tsx +++ b/src/app/org/[githubLogin]/_components/SidebarChat.tsx @@ -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"; @@ -727,8 +729,12 @@ function MessageArtifacts({ artifactIds }: { artifactIds?: string[] }) { return (
{artifacts.map((artifact) => { + if (artifact.type === "VERIFY" && isAuditVerdict(artifact.data)) { + return ( + + ); + } // Unknown artifact type — render nothing rather than crash. - void artifact; return null; })}
diff --git a/src/app/org/[githubLogin]/_state/useSendCanvasChatMessage.ts b/src/app/org/[githubLogin]/_state/useSendCanvasChatMessage.ts index 39a5eef73a..1c8a57bc52 100644 --- a/src/app/org/[githubLogin]/_state/useSendCanvasChatMessage.ts +++ b/src/app/org/[githubLogin]/_state/useSendCanvasChatMessage.ts @@ -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); diff --git a/src/app/w/[slug]/task/[...taskParams]/artifacts/index.ts b/src/app/w/[slug]/task/[...taskParams]/artifacts/index.ts index b88919e9e6..5ec71e1276 100644 --- a/src/app/w/[slug]/task/[...taskParams]/artifacts/index.ts +++ b/src/app/w/[slug]/task/[...taskParams]/artifacts/index.ts @@ -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"; diff --git a/src/app/w/[slug]/task/[...taskParams]/artifacts/verdict.tsx b/src/app/w/[slug]/task/[...taskParams]/artifacts/verdict.tsx new file mode 100644 index 0000000000..7a961ade8c --- /dev/null +++ b/src/app/w/[slug]/task/[...taskParams]/artifacts/verdict.tsx @@ -0,0 +1,226 @@ +"use client"; + +import { useState } from "react"; +import { Card } from "@/components/ui/card"; +import { Badge } from "@/components/ui/badge"; +import { Button } from "@/components/ui/button"; +import { Collapsible, CollapsibleContent, CollapsibleTrigger } from "@/components/ui/collapsible"; +import { Dialog, DialogContent, DialogHeader, DialogTitle } from "@/components/ui/dialog"; +import { Artifact, VerifyContent, VerifyOutcome } from "@/lib/chat"; +import { + ShieldCheck, + XCircle, + HelpCircle, + ChevronDown, + Camera, + Globe, + Terminal, + Clock, + FileText, + StickyNote, + Network, + Bug, + Database, +} from "lucide-react"; + +type BadgeVariant = "default" | "destructive" | "secondary"; + +const OUTCOME: Record< + VerifyOutcome, + { icon: typeof ShieldCheck; label: string; color: string; bg: string; border: string; badge: BadgeVariant } +> = { + works: { icon: ShieldCheck, label: "Verified", color: "text-green-500", bg: "bg-green-500/10", border: "border-green-500/30", badge: "default" }, + broken: { icon: XCircle, label: "Broken", color: "text-red-500", bg: "bg-red-500/10", border: "border-red-500/30", badge: "destructive" }, + unknown: { icon: HelpCircle, label: "Inconclusive", color: "text-amber-500", bg: "bg-amber-500/10", border: "border-amber-500/30", badge: "secondary" }, +}; + +const KIND_ICON: Record = { + screenshot: Camera, + http: Globe, + log: Terminal, + timing: Clock, + dom: FileText, + network: Network, + console: Bug, + db: Database, + note: StickyNote, +}; + +export function isAuditVerdict(content: unknown): content is VerifyContent { + return ( + !!content && + typeof content === "object" && + "overall" in content && + "claims" in content + ); +} + +export function VerdictPill({ artifact }: { artifact: Artifact }) { + const content = artifact.content as VerifyContent; + const [open, setOpen] = useState(false); + const outcome = OUTCOME[content.overall] ?? OUTCOME.unknown; + const Icon = outcome.icon; + const evidenceCount = content.evidence?.length ?? 0; + return ( + <> + + + + + + + Audit — {outcome.label} + + + + + + + ); +} + +export function VerdictArtifact({ artifact }: { artifact: Artifact }) { + const content = artifact.content as VerifyContent; + const [claimsOpen, setClaimsOpen] = useState(false); + const [detailsOpen, setDetailsOpen] = useState(false); + + const outcome = OUTCOME[content.overall] ?? OUTCOME.unknown; + const Icon = outcome.icon; + const claims = content.claims ?? []; + const evidence = content.evidence ?? []; + const observations = content.observations ?? []; + + const tally: Record = { works: 0, broken: 0, unknown: 0 }; + claims.forEach((c) => { + tally[c.verdict] = (tally[c.verdict] ?? 0) + 1; + }); + + return ( + +
+
+ +
+
+
+ {outcome.label} + Audit +
+

{content.summary}

+
+ {tally.works > 0 && {tally.works} works} + {tally.broken > 0 && {tally.broken} broken} + {tally.unknown > 0 && {tally.unknown} unknown} + {evidence.length > 0 && ( + · {evidence.length} evidence + )} +
+
+
+ + {claims.length > 0 && ( + + + + {claimsOpen ? "Hide" : "Show"} {claims.length} claim{claims.length === 1 ? "" : "s"} + + + + +
+ {claims.map((c, i) => { + const cfg = OUTCOME[c.verdict] ?? OUTCOME.unknown; + return ( +
+
+ {c.claim} + + {c.verdict} + +
+

{c.reasoning}

+ {c.proof?.length > 0 && ( +
+ {c.proof.map((p) => ( + + {p} + + ))} +
+ )} +
+ ); + })} +
+ {observations.length > 0 && ( +
+
Observations
+
    + {observations.map((o, i) => ( +
  • + {o} +
  • + ))} +
+
+ )} +
+
+ )} + + {evidence.length > 0 && ( +
+ +
+ )} + + + + + + + Audit — {outcome.label} + + +
+ {evidence.map((e) => { + const KindIcon = KIND_ICON[e.kind] ?? FileText; + return ( +
+
+ + {e.id} + {e.kind} + {e.summary} +
+ {e.kind === "screenshot" ? ( + {e.summary} + ) : e.data ? ( +
+                      {e.data}
+                    
+ ) : null} +
+ ); + })} +
+
+
+
+ ); +} diff --git a/src/app/w/[slug]/task/[...taskParams]/components/ChatMessage.tsx b/src/app/w/[slug]/task/[...taskParams]/components/ChatMessage.tsx index ccc8d07712..ea480d452f 100644 --- a/src/app/w/[slug]/task/[...taskParams]/components/ChatMessage.tsx +++ b/src/app/w/[slug]/task/[...taskParams]/components/ChatMessage.tsx @@ -7,7 +7,7 @@ import { ChevronDown, ChevronRight, Download, ExternalLink, User, X, Image as Im import { Button } from "@/components/ui/button"; import { ChatMessage as ChatMessageType, Option, FormContent } from "@/lib/chat"; -import { FormArtifact, LongformArtifactPanel, PublishWorkflowArtifact, PublishScriptArtifact, PublishPromptArtifact, PublishSkillArtifact, BountyArtifact } from "../artifacts"; +import { FormArtifact, LongformArtifactPanel, PublishWorkflowArtifact, PublishScriptArtifact, PublishPromptArtifact, PublishSkillArtifact, BountyArtifact, VerdictArtifact, isAuditVerdict } from "../artifacts"; import { PullRequestArtifact } from "../artifacts/pull-request"; import { MarkdownRenderer } from "@/components/MarkdownRenderer"; import { WorkflowUrlLink } from "./WorkflowUrlLink"; @@ -461,6 +461,17 @@ export const ChatMessage = memo(function ChatMessage({ ))} + {message.artifacts + ?.filter((a) => a.type === "VERIFY" && isAuditVerdict(a.content)) + .map((artifact) => ( +
+
+ + + +
+
+ ))} {message.artifacts ?.filter((a) => a.type === "PULL_REQUEST") .map((artifact) => ( diff --git a/src/config/middleware.ts b/src/config/middleware.ts index 0a49c1a819..bfc72fef44 100644 --- a/src/config/middleware.ts +++ b/src/config/middleware.ts @@ -169,6 +169,7 @@ export const ROUTE_POLICIES: ReadonlyArray = [ { path: "/api/ec2/alerts", strategy: "prefix", access: "webhook" }, { path: "/api/tasks/*/title", strategy: "pattern", access: "webhook" }, { path: "/api/tasks/*/recording", strategy: "pattern", access: "webhook" }, + { path: "/api/tasks/*/audit/callback", strategy: "pattern", access: "webhook" }, { path: "/api/tasks/*/webhook", strategy: "pattern", access: "webhook" }, { path: "/api/webhook/pool-manager", strategy: "prefix", access: "webhook" }, { path: "/api/w/*/pool/workspaces", strategy: "pattern", access: "webhook" }, diff --git a/src/lib/chat.ts b/src/lib/chat.ts index 99865009c5..b3d6fc78a4 100644 --- a/src/lib/chat.ts +++ b/src/lib/chat.ts @@ -59,6 +59,33 @@ export interface LongformContent { title?: string; } +export type VerifyOutcome = "works" | "broken" | "unknown"; + +export interface VerifyEvidence { + id: string; + kind: "screenshot" | "http" | "log" | "timing" | "dom" | "network" | "console" | "db" | "note"; + summary: string; + data: string; +} + +export interface VerifyClaim { + claim: string; + verdict: VerifyOutcome; + proof: string[]; + reasoning: string; +} + +export interface VerifyContent { + overall: VerifyOutcome; + claims: VerifyClaim[]; + observations: string[]; + summary: string; + evidence: VerifyEvidence[]; + startedAt?: string; + finishedAt?: string; + error?: string; +} + export interface BugReportContent { bugDescription: string; iframeUrl: string; @@ -332,7 +359,8 @@ export interface Artifact extends Omit { | PublishPromptContent | PublishSkillContent | BountyContent - | StreamContent; + | StreamContent + | VerifyContent; } // Using Prisma Attachment type directly (no additional fields needed) @@ -425,7 +453,8 @@ export function createArtifact(data: { | PublishPromptContent | PublishSkillContent | BountyContent - | StreamContent; + | StreamContent + | VerifyContent; icon?: ArtifactIcon; }): Artifact { return { diff --git a/src/lib/encryption/field-encryption.ts b/src/lib/encryption/field-encryption.ts index 251d1ef3f5..79b6f0bfca 100644 --- a/src/lib/encryption/field-encryption.ts +++ b/src/lib/encryption/field-encryption.ts @@ -68,6 +68,7 @@ export class FieldEncryptionService { "agentPassword", "agentWebhookSecret", "codeChangeWebhookSecret", + "auditCallbackKey", "vercelApiToken", "fiatPaymentPassword", "bifrostAdminPassword", diff --git a/src/lib/streaming/useStreamProcessor.ts b/src/lib/streaming/useStreamProcessor.ts index 49cff44334..88e1a19ed2 100644 --- a/src/lib/streaming/useStreamProcessor.ts +++ b/src/lib/streaming/useStreamProcessor.ts @@ -74,6 +74,7 @@ export function useStreamProcessor = []; // Unified timeline let error: string | undefined; let capturedUsage: TokenUsage | undefined; + const capturedArtifacts: Array<{ id: string; type: string; data: unknown }> = []; // Track text part sequence to generate unique IDs when stream reuses IDs let textPartSequence = 0; @@ -139,6 +140,7 @@ export function useStreamProcessor x.id === a.id)) { + capturedArtifacts.push(a); + debouncedUpdate(); + } } else if (data.type === "finish") { // Capture aggregated per-turn token usage from the AI SDK finish event. // Delegate field-name normalization to the shared utility so this diff --git a/src/services/auditor/deck.ts b/src/services/auditor/deck.ts new file mode 100644 index 0000000000..ce661ce873 --- /dev/null +++ b/src/services/auditor/deck.ts @@ -0,0 +1,159 @@ +import { db } from "@/lib/db"; +import { ArtifactType, ChatRole } from "@prisma/client"; +import type { Deck } from "./types"; + +export interface BuildDeckPod { + controlUrl?: string; + password?: string; + appUrl?: string; +} + +function resolveAppUrl(pod?: BuildDeckPod): string { + return ( + process.env.VERIFY_APP_URL || + pod?.appUrl || + process.env.NEXTAUTH_URL || + "http://localhost:3000" + ); +} + +function renderFeatureContext(feature: { + title: string; + brief: string | null; + requirements: string | null; + architecture: string | null; + personas: string[]; + userStories: { title: string }[]; +}): string { + const sections: string[] = [ + "=== CONTEXT ONLY — background for understanding the task, NOT part of what is being audited ===", + `Feature: ${feature.title}`, + ]; + + if (feature.brief) sections.push(`Brief:\n${feature.brief}`); + if (feature.requirements) sections.push(`Requirements:\n${feature.requirements}`); + if (feature.architecture) sections.push(`Architecture:\n${feature.architecture}`); + if (feature.personas.length > 0) sections.push(`Personas:\n${feature.personas.map((p) => `- ${p}`).join("\n")}`); + if (feature.userStories.length > 0) { + sections.push(`User Stories:\n${feature.userStories.map((s) => `- ${s.title}`).join("\n")}`); + } + + sections.push("=== END CONTEXT ONLY ==="); + return sections.join("\n\n"); +} + +async function loadEarliestUserAsks(taskId: string): Promise { + const messages = await db.chatMessage.findMany({ + where: { taskId, role: ChatRole.USER }, + orderBy: { createdAt: "asc" }, + take: 2, + select: { message: true }, + }); + + return messages + .map((m) => (m.message ?? "").trim()) + .filter((m) => m.length > 0); +} + +function assembleTaskPrompt(title: string, description: string | null, asks: string[]): string { + const parts: string[] = [title.trim()]; + + const trimmedDescription = description?.trim(); + if (trimmedDescription) { + parts.push(`Description:\n${trimmedDescription}`); + } + + if (asks.length > 0) { + parts.push(`Original request:\n${asks.join("\n\n")}`); + } + + return parts.join("\n\n"); +} + +async function loadDiffFromArtifacts(taskId: string): Promise { + const message = await db.chatMessage.findFirst({ + where: { + taskId, + artifacts: { some: { type: ArtifactType.DIFF } }, + }, + orderBy: { createdAt: "desc" }, + select: { + artifacts: { + where: { type: ArtifactType.DIFF }, + select: { content: true }, + }, + }, + }); + + const content = message?.artifacts[0]?.content; + if (!content) return null; + + return JSON.stringify(content); +} + +async function loadDiffFromPod(pod?: BuildDeckPod): Promise { + if (!pod?.controlUrl) return null; + + const headers: Record = {}; + if (pod.password) headers.Authorization = `Bearer ${pod.password}`; + + try { + const response = await fetch(`${pod.controlUrl}/diff`, { method: "GET", headers }); + if (!response.ok) { + console.error(`[auditor:deck] Failed to fetch diff from pod: ${response.status}`); + return null; + } + const diffs = await response.json(); + return JSON.stringify(diffs); + } catch (error) { + console.error("[auditor:deck] Error fetching diff from pod:", error); + return null; + } +} + +export async function buildDeck(taskId: string, pod?: BuildDeckPod): Promise { + const task = await db.task.findUnique({ + where: { id: taskId, deleted: false }, + select: { + title: true, + description: true, + feature: { + select: { + title: true, + brief: true, + requirements: true, + architecture: true, + personas: true, + userStories: { + orderBy: { order: "asc" }, + select: { title: true }, + }, + }, + }, + }, + }); + + if (!task) { + throw new Error(`Task not found: ${taskId}`); + } + + const diff = (await loadDiffFromArtifacts(taskId)) ?? (await loadDiffFromPod(pod)) ?? ""; + + const featureContext = task.feature ? renderFeatureContext(task.feature) : ""; + + const asks = await loadEarliestUserAsks(taskId); + const prompt = assembleTaskPrompt(task.title, task.description, asks); + + return { + task: { + prompt, + description: task.description ?? "", + }, + diff, + featureContext, + map: { + appUrl: resolveAppUrl(pod), + notes: null, + }, + }; +} diff --git a/src/services/auditor/trigger.ts b/src/services/auditor/trigger.ts new file mode 100644 index 0000000000..eb59d9e19e --- /dev/null +++ b/src/services/auditor/trigger.ts @@ -0,0 +1,158 @@ +import crypto from "node:crypto"; +import { db } from "@/lib/db"; +import { EncryptionService } from "@/lib/encryption"; +import { claimPodAndGetFrontend, updatePodRepositories, POD_PORTS } from "@/lib/pods"; +import { WorkflowStatus } from "@prisma/client"; +import { buildDeck } from "./deck"; +import { getDefaultModel, getApiKeyForModel } from "@/lib/ai/models"; +import type { AuditModel, AuditJobBody, StartAuditResult } from "./types"; + +const encryptionService = EncryptionService.getInstance(); + +async function buildModel(): Promise { + const registryModel = process.env.AUDIT_MODEL || (await getDefaultModel("task")); + if (registryModel) { + const apiKey = getApiKeyForModel(registryModel); + if (apiKey) { + const provider = registryModel.includes("/") ? registryModel.split("/")[0] : undefined; + return { apiKey, provider, model: registryModel }; + } + } + + if (process.env.ANTHROPIC_API_KEY) { + return { + apiKey: process.env.ANTHROPIC_API_KEY, + provider: "anthropic", + model: "anthropic/claude-opus-5", + }; + } + + return { + apiKey: process.env.OPENROUTER_API_KEY ?? "", + provider: "openrouter", + model: "openrouter/anthropic/claude-opus-4.8", + }; +} + +export async function startAudit(taskId: string, userId: string): Promise { + const task = await db.task.findUnique({ + where: { id: taskId, deleted: false }, + include: { + workspace: { + include: { + repositories: true, + swarm: true, + }, + }, + }, + }); + + if (!task) { + return { success: false, taskId, error: "Task not found" }; + } + + const customStakLinkUrl = process.env.CUSTOM_STAKLINK_URL; + + let controlUrl: string; + let podPassword = ""; + let podId: string | undefined; + let frontendUrl: string | undefined; + + if (customStakLinkUrl) { + controlUrl = customStakLinkUrl; + } else { + const swarm = task.workspace.swarm; + if (!swarm?.id || !swarm.poolApiKey) { + return { success: false, taskId, error: "Swarm not configured for pods" }; + } + + const poolApiKeyPlain = encryptionService.decryptField("poolApiKey", swarm.poolApiKey); + const services = swarm.services as Array<{ name: string; port: number }> | null | undefined; + + const podResult = await claimPodAndGetFrontend(swarm.id, poolApiKeyPlain, services || undefined); + + controlUrl = podResult.workspace.portMappings[POD_PORTS.CONTROL]; + podPassword = podResult.workspace.password; + podId = podResult.workspace.id; + frontendUrl = podResult.frontend; + + if (!controlUrl) { + return { success: false, taskId, error: "Control port not available on claimed pod" }; + } + + const repositories = task.workspace.repositories.map((r) => ({ url: r.repositoryUrl })); + if (repositories.length > 0) { + try { + await updatePodRepositories(controlUrl, podPassword, repositories); + } catch (error) { + console.error("[auditor:trigger] Failed to update pod repositories (non-fatal):", error); + } + } + } + + const deck = await buildDeck(taskId, { + controlUrl, + password: podPassword, + appUrl: frontendUrl, + }); + + const model = await buildModel(); + + const callbackApiKey = crypto.randomBytes(32).toString("hex"); + const encryptedCallbackKey = encryptionService.encryptField("auditCallbackKey", callbackApiKey); + + await db.task.update({ + where: { id: taskId }, + data: { + auditCallbackKey: JSON.stringify(encryptedCallbackKey), + ...(podId ? { podId } : {}), + }, + }); + + const baseUrl = process.env.NEXTAUTH_URL || "http://localhost:3000"; + const responseUrl = `${baseUrl}/api/tasks/${taskId}/audit/callback`; + + const body: AuditJobBody = { + taskId, + deck, + model, + responseUrl, + callbackApiKey, + }; + + const headers: Record = { "Content-Type": "application/json" }; + if (podPassword) headers.Authorization = `Bearer ${podPassword}`; + + try { + const response = await fetch(`${controlUrl}/audit`, { + method: "POST", + headers, + body: JSON.stringify(body), + }); + + if (!response.ok) { + const errorText = await response.text(); + throw new Error(`Pod returned ${response.status}: ${errorText}`); + } + } catch (error) { + console.error("[auditor:trigger] Failed to dispatch audit:", error); + await db.task.update({ + where: { id: taskId }, + data: { auditCallbackKey: null }, + }); + return { + success: false, + taskId, + error: error instanceof Error ? error.message : "Failed to dispatch audit", + }; + } + + await db.task.update({ + where: { id: taskId }, + data: { workflowStatus: WorkflowStatus.IN_PROGRESS, workflowStartedAt: new Date() }, + }); + + console.log(`[auditor:trigger] Audit dispatched for task ${taskId} by user ${userId}`); + + return { success: true, taskId, podId, appUrl: deck.map.appUrl }; +} diff --git a/src/services/auditor/types.ts b/src/services/auditor/types.ts new file mode 100644 index 0000000000..a833cf9bf0 --- /dev/null +++ b/src/services/auditor/types.ts @@ -0,0 +1,69 @@ +export interface DeckTask { + prompt: string; + description: string; +} + +export interface DeckMap { + appUrl: string; + notes: string | null; +} + +export interface Deck { + task: DeckTask; + diff: string; + featureContext: string; + map: DeckMap; +} + +export interface AuditModel { + apiKey: string; + host?: string; + provider?: string; + model?: string; +} + +export interface AuditJobBody { + taskId: string; + deck: Deck; + model: AuditModel; + responseUrl: string; + callbackApiKey: string; +} + +export type AuditOverall = "works" | "broken" | "unknown"; + +export interface AuditClaim { + claim: string; + verdict: string; + proof: string[]; + reasoning: string; +} + +export type AuditEvidenceKind = "screenshot" | "http" | "log" | "timing" | "dom" | "note"; + +export interface AuditEvidence { + id: string; + kind: AuditEvidenceKind; + summary: string; + data: string; +} + +export interface AuditVerdict { + taskId: string; + overall: AuditOverall; + claims: AuditClaim[]; + observations: string[]; + summary: string; + evidence?: AuditEvidence[]; + startedAt: string; + finishedAt: string; + error?: string; +} + +export interface StartAuditResult { + success: boolean; + taskId: string; + podId?: string; + appUrl?: string; + error?: string; +} diff --git a/src/types/encryption.ts b/src/types/encryption.ts index 1201181cd6..2a487e1b6d 100644 --- a/src/types/encryption.ts +++ b/src/types/encryption.ts @@ -31,6 +31,7 @@ export type EncryptableField = | "agentPassword" | "agentWebhookSecret" | "codeChangeWebhookSecret" + | "auditCallbackKey" | "vercelApiToken" | "vercelWebhookSecret" | "sphinxBotSecret" diff --git a/src/types/streaming.ts b/src/types/streaming.ts index 48b8bb7769..a752aa37b0 100644 --- a/src/types/streaming.ts +++ b/src/types/streaming.ts @@ -18,6 +18,7 @@ export type StreamEventType = | "start" | "finish" | "data-usage" + | "data-artifact" | "error"; export type ToolCallStatus = @@ -152,6 +153,12 @@ export interface UsageUpdateEvent extends BaseStreamEvent { data: TokenUsage; } +/** A chat artifact (e.g. a verification verdict) emitted mid-stream by the server. */ +export interface ArtifactEvent extends BaseStreamEvent { + type: "data-artifact"; + data: { id: string; type: string; data: unknown }; +} + export type StreamEvent = | TextStartEvent | TextDeltaEvent @@ -169,6 +176,7 @@ export type StreamEvent = | StartEvent | FinishEvent | UsageUpdateEvent + | ArtifactEvent | ErrorEvent; // Generic streaming message structure @@ -213,6 +221,8 @@ export interface BaseStreamingMessage { error?: string; /** Per-turn aggregated token usage, populated from the `finish` SSE event. */ usage?: TokenUsage; + /** Chat artifacts (e.g. verification verdicts) emitted mid-stream via `data-artifact`. */ + artifacts?: Array<{ id: string; type: string; data: unknown }>; } // Tool processor function type