From ecea8655dee94b2acb1d1343a281f836613bd3c3 Mon Sep 17 00:00:00 2001 From: ashishsinghbora <135435891+ashishsinghbora@users.noreply.github.com> Date: Sat, 26 Sep 2026 00:44:48 +0530 Subject: [PATCH] feat(sprint-3-4): PR lifecycle, issue linking, review ingestion, webhook idempotency, and repo trust boundary --- drizzle/0003_pr_nullable_user.sql | 1 + src/app/api/webhooks/github/route.ts | 153 +++------ src/lib/db/schema.ts | 1 - src/lib/github/webhooks.test.ts | 410 ++++++++++++++++++++-- src/lib/github/webhooks.ts | 496 ++++++++++++++++++++++++++- 5 files changed, 936 insertions(+), 125 deletions(-) create mode 100644 drizzle/0003_pr_nullable_user.sql diff --git a/drizzle/0003_pr_nullable_user.sql b/drizzle/0003_pr_nullable_user.sql new file mode 100644 index 0000000..a726880 --- /dev/null +++ b/drizzle/0003_pr_nullable_user.sql @@ -0,0 +1 @@ +ALTER TABLE "pull_requests" ALTER COLUMN "user_id" DROP NOT NULL; diff --git a/src/app/api/webhooks/github/route.ts b/src/app/api/webhooks/github/route.ts index dba469f..c28b18c 100644 --- a/src/app/api/webhooks/github/route.ts +++ b/src/app/api/webhooks/github/route.ts @@ -1,12 +1,20 @@ import { NextRequest, NextResponse } from "next/server"; -import { verifyWebhookSignature, type WebhookPullRequestEvent } from "@/lib/github/webhooks"; -import { getDb, schema } from "@/lib/db"; -import { eq, and } from "drizzle-orm"; -import { processProgressionOnContribution } from "@/lib/progression/pipeline"; +import { + verifyWebhookSignature, + handlePullRequestWebhook, + handlePullRequestReviewWebhook, + isWebhookDeliveryProcessed, + recordWebhookDelivery, + markWebhookDeliveryCompleted, + markWebhookDeliveryFailed, + type WebhookPullRequestEvent, + type WebhookPullRequestReviewEvent, +} from "@/lib/github/webhooks"; export async function POST(request: NextRequest) { const signature = request.headers.get("x-hub-signature-256"); const event = request.headers.get("x-github-event"); + const deliveryId = request.headers.get("x-github-delivery"); const rawBody = await request.text(); // Validate webhook secret signature @@ -18,109 +26,52 @@ export async function POST(request: NextRequest) { return NextResponse.json({ message: "PONG" }); } - if (event === "pull_request") { - try { - const payload = JSON.parse(rawBody) as WebhookPullRequestEvent; - - if (payload.action === "closed" && payload.pull_request.merged) { - const db = await getDb(); - const pr = payload.pull_request; - const repoFullName = payload.repository.full_name; - - // 1. Look up user in database - const matchingUsers = await db - .select() - .from(schema.users) - .where(eq(schema.users.githubId, pr.user.id)) - .limit(1); - - if (matchingUsers.length === 0) { - return NextResponse.json({ message: "PR author not registered on platform, skipped" }); - } - - const user = matchingUsers[0]; - - // 2. Find or create project entry - const matchingProjects = await db - .select() - .from(schema.projects) - .where(eq(schema.projects.githubRepo, repoFullName)) - .limit(1); + // Idempotency: Ignore duplicate webhook deliveries + if (deliveryId) { + const alreadyProcessed = await isWebhookDeliveryProcessed(deliveryId); + if (alreadyProcessed) { + return NextResponse.json({ + message: "Duplicate webhook delivery ignored", + deliveryId, + }); + } - let projectId: string; - if (matchingProjects.length === 0) { - projectId = `proj_${crypto.randomUUID()}`; - await db.insert(schema.projects).values({ - id: projectId, - name: payload.repository.name, - slug: payload.repository.name.toLowerCase(), - githubRepo: repoFullName, - description: `Official repository: ${repoFullName}`, - primaryLanguage: "TypeScript", - languages: [], - isOfficial: true, - }); - } else { - projectId = matchingProjects[0].id; - } + let repoName: string | undefined; + try { + const parsed = JSON.parse(rawBody); + repoName = parsed.repository?.full_name; + } catch { + // ignore + } - // 3. Check existing merged contributions for this user - const existingContributions = await db - .select() - .from(schema.contributions) - .where( - and( - eq(schema.contributions.userId, user.id), - eq(schema.contributions.state, "merged") - ) - ); + await recordWebhookDelivery({ + deliveryId, + eventType: event || "unknown", + repository: repoName, + }); + } - const isFirstPr = existingContributions.length === 0; - const contributionId = `contrib_${crypto.randomUUID()}`; - const mergedDate = pr.merged_at ? new Date(pr.merged_at) : new Date(); + try { + let result: any = { received: true, event }; - await db.insert(schema.contributions).values({ - id: contributionId, - userId: user.id, - projectId, - githubPrNumber: pr.number, - prTitle: pr.title, - prUrl: pr.html_url, - state: "merged", - isFirstPr, - mergedAt: mergedDate, - verifiedAt: new Date(), - verificationSource: "github_webhook", - }); + if (event === "pull_request") { + const payload = JSON.parse(rawBody) as WebhookPullRequestEvent; + result = await handlePullRequestWebhook(payload); + } else if (event === "pull_request_review") { + const payload = JSON.parse(rawBody) as WebhookPullRequestReviewEvent; + result = await handlePullRequestReviewWebhook(payload); + } - // 4. Process contributor progression, credentials, and Founding 1,000 cohort - const progressionResult = await processProgressionOnContribution(db, { - userId: user.id, - contributionId, - githubUsername: user.githubUsername, - repoFullName, - prNumber: pr.number, - prUrl: pr.html_url, - prTitle: pr.title, - mergedAt: mergedDate, - isFirstPr, - }); + if (deliveryId) { + await markWebhookDeliveryCompleted(deliveryId); + } - return NextResponse.json({ - status: "success", - contributionId, - isFirstPr, - promoted: progressionResult.promoted, - level: progressionResult.newLevel, - credentialsIssued: progressionResult.credentialsIssued, - foundingMemberNumber: progressionResult.foundingMemberNumber, - }); - } - } catch (err: any) { - console.error("Webhook processing error:", err); - return NextResponse.json({ error: err.message }, { status: 500 }); + return NextResponse.json(result); + } catch (err: any) { + console.error("Webhook processing error:", err); + if (deliveryId) { + await markWebhookDeliveryFailed(deliveryId, err.message); } + return NextResponse.json({ error: err.message }, { status: 500 }); } - - return NextResponse.json({ received: true }); } diff --git a/src/lib/db/schema.ts b/src/lib/db/schema.ts index b217842..dffb821 100644 --- a/src/lib/db/schema.ts +++ b/src/lib/db/schema.ts @@ -345,7 +345,6 @@ export const pullRequests = pgTable( { id: text("id").primaryKey(), userId: text("user_id") - .notNull() .references(() => users.id, { onDelete: "cascade" }), projectId: text("project_id") .notNull() diff --git a/src/lib/github/webhooks.test.ts b/src/lib/github/webhooks.test.ts index 8a8fdae..15fa647 100644 --- a/src/lib/github/webhooks.test.ts +++ b/src/lib/github/webhooks.test.ts @@ -1,34 +1,400 @@ -import { describe, it, expect } from "vitest"; -import { verifyWebhookSignature } from "./webhooks"; +import { describe, it, expect, beforeAll } from "vitest"; +import { + verifyWebhookSignature, + extractIssueNumbersFromPrBody, + handlePullRequestWebhook, + handlePullRequestReviewWebhook, + recordWebhookDelivery, + isWebhookDeliveryProcessed, + markWebhookDeliveryCompleted, + type WebhookPullRequestEvent, + type WebhookPullRequestReviewEvent, +} from "./webhooks"; import crypto from "crypto"; +import { getDb, schema } from "@/lib/db"; +import { runMigrations } from "@/lib/db/migrate"; +import { eq } from "drizzle-orm"; -describe("GitHub Webhook Signature Verification", () => { +describe("GitHub Webhooks & PR Lifecycle Engine", () => { const secret = process.env.GITHUB_WEBHOOK_SECRET || "dev_webhook_secret_key"; - const payload = JSON.stringify({ action: "closed", pull_request: { merged: true } }); + const userId = "usr_wh_test_contributor"; + const projectId = "proj_wh_test_repo"; + const issueId = "iss_wh_test_issue"; + const claimId = "claim_wh_test_claim"; - it("verifies valid HMAC SHA256 signature", () => { - const signature = `sha256=${crypto - .createHmac("sha256", secret) - .update(payload) - .digest("hex")}`; + beforeAll(async () => { + (process.env as Record).NODE_ENV = "test"; + await runMigrations(); + const db = await getDb(); - const isValid = verifyWebhookSignature(payload, signature); - expect(isValid).toBe(true); + // 1. Seed user + await db.insert(schema.users).values({ + id: userId, + githubId: 998877, + githubUsername: "contributor_cathy", + displayName: "Cathy Contributor", + role: "contributor", + level: "level_1", + isOnboarded: true, + }); + + // 2. Seed project + await db.insert(schema.projects).values({ + id: projectId, + name: "awesome-sdk", + slug: "awesome-sdk", + githubRepo: "TechNexusOrg/awesome-sdk", + description: "Official test repository for webhook lifecycle", + primaryLanguage: "TypeScript", + contributionEnabled: true, + }); + + // 3. Seed issue + await db.insert(schema.issues).values({ + id: issueId, + projectId, + githubIssueId: 7771, + githubIssueNumber: 42, + title: "Add rate limiting retry strategy", + state: "open", + htmlUrl: "https://github.com/TechNexusOrg/awesome-sdk/issues/42", + difficulty: "beginner", + }); + + // 4. Seed active issue claim by user + await db.insert(schema.issueClaims).values({ + id: claimId, + issueId, + userId, + status: "active", + claimedAt: new Date(), + expiresAt: new Date(Date.now() + 7 * 24 * 60 * 60 * 1000), + }); }); - it("rejects forged or modified payload", () => { - const signature = `sha256=${crypto - .createHmac("sha256", secret) - .update(payload) - .digest("hex")}`; + describe("HMAC SHA256 Signature Verification", () => { + const payload = JSON.stringify({ action: "closed", pull_request: { merged: true } }); + + it("verifies valid HMAC SHA256 signature", () => { + const signature = `sha256=${crypto + .createHmac("sha256", secret) + .update(payload) + .digest("hex")}`; + + const isValid = verifyWebhookSignature(payload, signature); + expect(isValid).toBe(true); + }); - const modifiedPayload = JSON.stringify({ action: "closed", pull_request: { merged: false } }); - const isValid = verifyWebhookSignature(modifiedPayload, signature); - expect(isValid).toBe(false); + it("rejects forged or modified payload", () => { + const signature = `sha256=${crypto + .createHmac("sha256", secret) + .update(payload) + .digest("hex")}`; + + const modifiedPayload = JSON.stringify({ action: "closed", pull_request: { merged: false } }); + const isValid = verifyWebhookSignature(modifiedPayload, signature); + expect(isValid).toBe(false); + }); + + it("rejects missing or malformed signature header", () => { + expect(verifyWebhookSignature(payload, null)).toBe(false); + expect(verifyWebhookSignature(payload, "invalid_sig")).toBe(false); + }); }); - it("rejects missing or malformed signature header", () => { - expect(verifyWebhookSignature(payload, null)).toBe(false); - expect(verifyWebhookSignature(payload, "invalid_sig")).toBe(false); + describe("Issue Number Extraction from PR Body", () => { + it("extracts issue numbers from standard closing keywords", () => { + const body1 = "This PR implements exponential backoff. Fixes #42 cleanly."; + expect(extractIssueNumbersFromPrBody(body1)).toEqual([42]); + + const body2 = "Closes #15 and also resolves #89"; + expect(extractIssueNumbersFromPrBody(body2)).toEqual([15, 89]); + + const body3 = "Resolves https://github.com/TechNexusOrg/platform/issues/99"; + expect(extractIssueNumbersFromPrBody(body3)).toEqual([99]); + }); + + it("returns empty array for text without issue references", () => { + expect(extractIssueNumbersFromPrBody("Just refactoring variable names")).toEqual([]); + expect(extractIssueNumbersFromPrBody(null)).toEqual([]); + expect(extractIssueNumbersFromPrBody(undefined)).toEqual([]); + }); + }); + + describe("Pull Request Lifecycle & Issue Claim Auto-Completion", () => { + const prNumber = 105; + const prUrl = "https://github.com/TechNexusOrg/awesome-sdk/pull/105"; + + it("ingests opened PR, links to referenced issue and registered user", async () => { + const db = await getDb(); + const openEvent: WebhookPullRequestEvent = { + action: "opened", + pull_request: { + id: 55001, + number: prNumber, + title: "feat: add rate limiting retry", + body: "Implements retry strategy. Fixes #42", + html_url: prUrl, + state: "open", + draft: false, + merged: false, + merged_at: null, + user: { + id: 998877, + login: "contributor_cathy", + avatar_url: "https://avatars.githubusercontent.com/u/998877", + }, + base: { + repo: { + id: 33001, + name: "awesome-sdk", + full_name: "TechNexusOrg/awesome-sdk", + html_url: "https://github.com/TechNexusOrg/awesome-sdk", + }, + }, + }, + repository: { + id: 33001, + name: "awesome-sdk", + full_name: "TechNexusOrg/awesome-sdk", + owner: { + login: "TechNexusOrg", + }, + }, + sender: { + id: 998877, + login: "contributor_cathy", + }, + }; + + const result = await handlePullRequestWebhook(openEvent, db); + expect(result.handled).toBe(true); + expect(result.linkedIssueId).toBe(issueId); + + // Verify record in pull_requests table + const [storedPr] = await db + .select() + .from(schema.pullRequests) + .where(eq(schema.pullRequests.url, prUrl)) + .limit(1); + + expect(storedPr).toBeDefined(); + expect(storedPr.userId).toBe(userId); + expect(storedPr.issueId).toBe(issueId); + expect(storedPr.state).toBe("open"); + }); + + it("ingests merged PR, completes active claim, closes issue, and awards contribution", async () => { + const db = await getDb(); + const mergeEvent: WebhookPullRequestEvent = { + action: "closed", + pull_request: { + id: 55001, + number: prNumber, + title: "feat: add rate limiting retry", + body: "Implements retry strategy. Fixes #42", + html_url: prUrl, + state: "closed", + draft: false, + merged: true, + merged_at: new Date().toISOString(), + merge_commit_sha: "abcd1234efgh5678", + user: { + id: 998877, + login: "contributor_cathy", + avatar_url: "https://avatars.githubusercontent.com/u/998877", + }, + base: { + repo: { + id: 33001, + name: "awesome-sdk", + full_name: "TechNexusOrg/awesome-sdk", + html_url: "https://github.com/TechNexusOrg/awesome-sdk", + }, + }, + }, + repository: { + id: 33001, + name: "awesome-sdk", + full_name: "TechNexusOrg/awesome-sdk", + owner: { + login: "TechNexusOrg", + }, + }, + sender: { + id: 998877, + login: "contributor_cathy", + }, + }; + + const result = await handlePullRequestWebhook(mergeEvent, db); + expect(result.handled).toBe(true); + expect(result.contributionId).toBeDefined(); + + // 1. Verify PR state updated to 'merged' + const [updatedPr] = await db + .select() + .from(schema.pullRequests) + .where(eq(schema.pullRequests.url, prUrl)) + .limit(1); + expect(updatedPr.state).toBe("merged"); + + // 2. Verify issue state updated to 'closed' + const [closedIssue] = await db + .select() + .from(schema.issues) + .where(eq(schema.issues.id, issueId)) + .limit(1); + expect(closedIssue.state).toBe("closed"); + + // 3. Verify issue claim updated to 'completed' + const [completedClaim] = await db + .select() + .from(schema.issueClaims) + .where(eq(schema.issueClaims.id, claimId)) + .limit(1); + expect(completedClaim.status).toBe("completed"); + + // 4. Verify contribution recorded + const [contrib] = await db + .select() + .from(schema.contributions) + .where(eq(schema.contributions.prUrl, prUrl)) + .limit(1); + expect(contrib).toBeDefined(); + expect(contrib.state).toBe("merged"); + expect(contrib.issueId).toBe(issueId); + }); + }); + + describe("Pull Request Code Reviews", () => { + it("ingests approved review and updates PR state", async () => { + const db = await getDb(); + const prUrl = "https://github.com/TechNexusOrg/awesome-sdk/pull/105"; + + const reviewEvent: WebhookPullRequestReviewEvent = { + action: "submitted", + review: { + id: 881122, + state: "approved", + html_url: `${prUrl}#pullrequestreview-881122`, + submitted_at: new Date().toISOString(), + user: { + id: 112233, + login: "maintainer_dan", + avatar_url: "https://avatars.githubusercontent.com/u/112233", + }, + }, + pull_request: { + id: 55001, + number: 105, + title: "feat: add rate limiting retry", + html_url: prUrl, + state: "open", + }, + repository: { + id: 33001, + name: "awesome-sdk", + full_name: "TechNexusOrg/awesome-sdk", + }, + sender: { + id: 112233, + login: "maintainer_dan", + }, + }; + + const result = await handlePullRequestReviewWebhook(reviewEvent, db); + expect(result.handled).toBe(true); + expect(result.reviewState).toBe("approved"); + + // Verify review recorded in pull_request_reviews table + const [storedReview] = await db + .select() + .from(schema.pullRequestReviews) + .where(eq(schema.pullRequestReviews.githubReviewId, 881122)) + .limit(1); + + expect(storedReview).toBeDefined(); + expect(storedReview.reviewerUsername).toBe("maintainer_dan"); + expect(storedReview.reviewState).toBe("approved"); + }); + }); + + describe("Webhook Delivery Idempotency Tracking", () => { + const testDeliveryId = "delivery_test_unique_guid_123"; + + it("records delivery and marks it completed, detecting duplicates", async () => { + const db = await getDb(); + + // Check initial state + const initialProcessed = await isWebhookDeliveryProcessed(testDeliveryId, db); + expect(initialProcessed).toBe(false); + + // Record incoming delivery + await recordWebhookDelivery( + { + deliveryId: testDeliveryId, + eventType: "pull_request", + repository: "TechNexusOrg/awesome-sdk", + }, + db + ); + + // Mark completed + await markWebhookDeliveryCompleted(testDeliveryId, db); + + // Now it should be detected as processed + const isDone = await isWebhookDeliveryProcessed(testDeliveryId, db); + expect(isDone).toBe(true); + }); + }); + + describe("Repository Trust Boundary", () => { + it("rejects PR events originating from unauthorized third-party repositories", async () => { + const db = await getDb(); + const maliciousEvent: WebhookPullRequestEvent = { + action: "opened", + pull_request: { + id: 99999, + number: 1, + title: "Malicious external PR", + body: "Trying to claim contribution outside org", + html_url: "https://github.com/RandomSpammer/fake-repo/pull/1", + state: "open", + merged: false, + merged_at: null, + user: { + id: 998877, + login: "contributor_cathy", + avatar_url: "https://avatars.githubusercontent.com/u/998877", + }, + base: { + repo: { + id: 88888, + name: "fake-repo", + full_name: "RandomSpammer/fake-repo", + html_url: "https://github.com/RandomSpammer/fake-repo", + }, + }, + }, + repository: { + id: 88888, + name: "fake-repo", + full_name: "RandomSpammer/fake-repo", + owner: { + login: "RandomSpammer", + }, + }, + sender: { + id: 998877, + login: "contributor_cathy", + }, + }; + + const result = await handlePullRequestWebhook(maliciousEvent, db); + expect(result.handled).toBe(false); + expect(result.reason).toMatch(/outside official organization trust boundary/i); + }); }); }); + diff --git a/src/lib/github/webhooks.ts b/src/lib/github/webhooks.ts index 30f9538..8b89e8c 100644 --- a/src/lib/github/webhooks.ts +++ b/src/lib/github/webhooks.ts @@ -1,6 +1,12 @@ import crypto from "crypto"; import { env } from "@/lib/env"; +import { getDb, schema } from "@/lib/db"; +import { eq, and } from "drizzle-orm"; +import { processProgressionOnContribution } from "@/lib/progression/pipeline"; +/** + * Validates HMAC SHA256 webhook signature against configured secret using timingSafeEqual. + */ export function verifyWebhookSignature(payloadBody: string, signatureHeader: string | null): boolean { if (!signatureHeader || !signatureHeader.startsWith("sha256=")) { return false; @@ -26,16 +32,53 @@ export function verifyWebhookSignature(payloadBody: string, signatureHeader: str return crypto.timingSafeEqual(expectedBuffer, actualBuffer); } +/** + * Extracts referenced issue numbers from a pull request title and body using standard GitHub keywords. + * Examples: "Fixes #12", "Closes https://github.com/TechNexusOrg/platform/issues/45", "Resolves #8" + */ +export function extractIssueNumbersFromPrBody(body: string | null | undefined): number[] { + if (!body) return []; + + const issueNumbers = new Set(); + + // Match keyword closing syntax: Fixes #123, Closes #45, Resolves #89 + const keywordRegex = + /(?:close|closes|closed|fix|fixes|fixed|resolve|resolves|resolved)\s+(?:#|https?:\/\/github\.com\/[^\/\s]+\/[^\/\s]+\/issues\/)(\d+)/gi; + + let match: RegExpExecArray | null; + while ((match = keywordRegex.exec(body)) !== null) { + const num = parseInt(match[1], 10); + if (!isNaN(num)) { + issueNumbers.add(num); + } + } + + // Also match fallback "#" when explicitly written + const hashtagRegex = /#(\d+)/g; + while ((match = hashtagRegex.exec(body)) !== null) { + const num = parseInt(match[1], 10); + if (!isNaN(num)) { + issueNumbers.add(num); + } + } + + return Array.from(issueNumbers); +} + export interface WebhookPullRequestEvent { - action: "opened" | "closed" | "reopened" | "synchronize"; + action: "opened" | "closed" | "reopened" | "synchronize" | "edited"; pull_request: { id: number; number: number; title: string; + body?: string | null; html_url: string; state: "open" | "closed"; + draft?: boolean; merged: boolean; merged_at: string | null; + closed_at?: string | null; + merge_commit_sha?: string | null; user: { id: number; login: string; @@ -63,3 +106,454 @@ export interface WebhookPullRequestEvent { login: string; }; } + +export interface WebhookPullRequestReviewEvent { + action: "submitted" | "edited" | "dismissed"; + review: { + id: number; + state: string; // "approved" | "changes_requested" | "commented" | "dismissed" + html_url: string; + submitted_at: string; + user: { + id: number; + login: string; + avatar_url: string; + }; + }; + pull_request: { + id: number; + number: number; + title: string; + html_url: string; + state: "open" | "closed"; + }; + repository: { + id: number; + name: string; + full_name: string; + }; + sender: { + id: number; + login: string; + }; +} + +/** + * Checks if a GitHub webhook delivery ID has already been successfully processed. + */ +export async function isWebhookDeliveryProcessed( + deliveryId: string, + database?: any +): Promise { + const db = database || (await getDb()); + const existing = await db + .select({ status: schema.githubWebhookDeliveries.status }) + .from(schema.githubWebhookDeliveries) + .where(eq(schema.githubWebhookDeliveries.deliveryId, deliveryId)) + .limit(1); + + return existing.length > 0 && existing[0].status === "completed"; +} + +/** + * Records an incoming webhook delivery with initial status. + */ +export async function recordWebhookDelivery( + params: { + deliveryId: string; + eventType: string; + repository?: string; + }, + database?: any +) { + const db = database || (await getDb()); + const now = new Date(); + + const existing = await db + .select({ id: schema.githubWebhookDeliveries.id }) + .from(schema.githubWebhookDeliveries) + .where(eq(schema.githubWebhookDeliveries.deliveryId, params.deliveryId)) + .limit(1); + + if (existing.length === 0) { + await db.insert(schema.githubWebhookDeliveries).values({ + id: `whdel_${crypto.randomUUID()}`, + deliveryId: params.deliveryId, + eventType: params.eventType, + repository: params.repository || null, + status: "processing", + receivedAt: now, + }); + } +} + +/** + * Marks a webhook delivery as completed. + */ +export async function markWebhookDeliveryCompleted( + deliveryId: string, + database?: any +) { + const db = database || (await getDb()); + await db + .update(schema.githubWebhookDeliveries) + .set({ + status: "completed", + processedAt: new Date(), + }) + .where(eq(schema.githubWebhookDeliveries.deliveryId, deliveryId)); +} + +/** + * Marks a webhook delivery as failed. + */ +export async function markWebhookDeliveryFailed( + deliveryId: string, + error: string, + database?: any +) { + const db = database || (await getDb()); + await db + .update(schema.githubWebhookDeliveries) + .set({ + status: "failed", + error, + processedAt: new Date(), + }) + .where(eq(schema.githubWebhookDeliveries.deliveryId, deliveryId)); +} + +/** + * Handles pull_request webhook events: + * - Upserts PR record into `pull_requests` + * - Links PR to referenced issues and active claims + * - On merge: auto-completes claims, closes issues, records verified contribution, updates progression + */ +export async function handlePullRequestWebhook( + payload: WebhookPullRequestEvent, + database?: any +) { + const db = database || (await getDb()); + const pr = payload.pull_request; + const repoFullName = payload.repository.full_name; + const now = new Date(); + + // 1. Resolve or create project with repository trust boundary enforcement + const matchingProjects = await db + .select() + .from(schema.projects) + .where(eq(schema.projects.githubRepo, repoFullName)) + .limit(1); + + let projectId: string; + let isContributionEligible = false; + + if (matchingProjects.length === 0) { + const isTechNexusOrg = repoFullName.toLowerCase().startsWith("technexusorg/"); + if (!isTechNexusOrg) { + return { + handled: false, + reason: "Repository outside official organization trust boundary.", + prId: null, + }; + } + + projectId = `proj_${crypto.randomUUID()}`; + await db.insert(schema.projects).values({ + id: projectId, + name: payload.repository.name, + slug: payload.repository.name.toLowerCase(), + githubRepo: repoFullName, + description: `Official repository: ${repoFullName}`, + primaryLanguage: "TypeScript", + languages: [], + isOfficial: true, + contributionEnabled: true, + }); + isContributionEligible = true; + } else { + const project = matchingProjects[0]; + projectId = project.id; + isContributionEligible = project.isOfficial && project.contributionEnabled; + } + + if (!isContributionEligible) { + return { + handled: false, + reason: "Repository is outside official contribution scope or contributions are disabled.", + prId: null, + }; + } + + // 2. Resolve user if registered on platform + const matchingUsers = await db + .select() + .from(schema.users) + .where(eq(schema.users.githubId, pr.user.id)) + .limit(1); + + const registeredUser = matchingUsers.length > 0 ? matchingUsers[0] : null; + + // 3. Extract and match referenced issues + const fullText = `${pr.title} ${pr.body || ""}`; + const issueNumbers = extractIssueNumbersFromPrBody(fullText); + + let matchedIssueId: string | null = null; + if (issueNumbers.length > 0) { + // Check if any referenced issue exists for this project in our database + for (const num of issueNumbers) { + const issues = await db + .select({ id: schema.issues.id }) + .from(schema.issues) + .where( + and( + eq(schema.issues.projectId, projectId), + eq(schema.issues.githubIssueNumber, num) + ) + ) + .limit(1); + + if (issues.length > 0) { + matchedIssueId = issues[0].id; + break; + } + } + } + + // 4. Map PR state + let prState: "open" | "approved" | "changes_requested" | "merged" | "closed" = "open"; + if (pr.merged) { + prState = "merged"; + } else if (pr.state === "closed") { + prState = "closed"; + } + + // 5. Upsert into pull_requests table + const existingPrs = await db + .select() + .from(schema.pullRequests) + .where(eq(schema.pullRequests.url, pr.html_url)) + .limit(1); + + let prRecordId: string; + if (existingPrs.length > 0) { + prRecordId = existingPrs[0].id; + await db + .update(schema.pullRequests) + .set({ + title: pr.title, + state: prState, + draft: pr.draft || false, + mergeCommitSha: pr.merge_commit_sha || null, + issueId: matchedIssueId || existingPrs[0].issueId, + mergedAt: pr.merged_at ? new Date(pr.merged_at) : null, + closedAt: pr.closed_at ? new Date(pr.closed_at) : null, + updatedAt: now, + }) + .where(eq(schema.pullRequests.id, prRecordId)); + } else { + prRecordId = `pr_${crypto.randomUUID()}`; + await db.insert(schema.pullRequests).values({ + id: prRecordId, + userId: registeredUser?.id || null, + projectId, + issueId: matchedIssueId, + githubPrId: pr.id, + githubPrNumber: pr.number, + title: pr.title, + url: pr.html_url, + state: prState, + draft: pr.draft || false, + mergeCommitSha: pr.merge_commit_sha || null, + mergedAt: pr.merged_at ? new Date(pr.merged_at) : null, + closedAt: pr.closed_at ? new Date(pr.closed_at) : null, + createdAt: now, + updatedAt: now, + }); + } + + // 6. If PR is merged: close issue, complete active claims, and record verified contribution + let contributionId: string | undefined; + let progressionResult: any | undefined; + + if (payload.action === "closed" && pr.merged) { + // A. If an issue was linked, complete the active claim and close the issue + if (matchedIssueId) { + // Mark issue closed + await db + .update(schema.issues) + .set({ + state: "closed", + updatedAt: now, + }) + .where(eq(schema.issues.id, matchedIssueId)); + + // Mark active claim completed + await db + .update(schema.issueClaims) + .set({ + status: "completed", + updatedAt: now, + }) + .where( + and( + eq(schema.issueClaims.issueId, matchedIssueId), + eq(schema.issueClaims.status, "active") + ) + ); + } + + // B. If the author is a registered contributor, record verified contribution & evaluate progression + if (registeredUser) { + const existingContribs = await db + .select() + .from(schema.contributions) + .where( + and( + eq(schema.contributions.userId, registeredUser.id), + eq(schema.contributions.state, "merged") + ) + ); + + const isFirstPr = existingContribs.length === 0; + contributionId = `contrib_${crypto.randomUUID()}`; + const mergedDate = pr.merged_at ? new Date(pr.merged_at) : now; + + // Upsert contribution by prUrl + const existingByUrl = await db + .select() + .from(schema.contributions) + .where(eq(schema.contributions.prUrl, pr.html_url)) + .limit(1); + + if (existingByUrl.length === 0) { + await db.insert(schema.contributions).values({ + id: contributionId, + userId: registeredUser.id, + projectId, + issueId: matchedIssueId, + githubPrNumber: pr.number, + prTitle: pr.title, + prUrl: pr.html_url, + state: "merged", + isFirstPr, + mergedAt: mergedDate, + verifiedAt: now, + verificationSource: "github_webhook", + }); + + progressionResult = await processProgressionOnContribution(db, { + userId: registeredUser.id, + contributionId, + githubUsername: registeredUser.githubUsername, + repoFullName, + prNumber: pr.number, + prUrl: pr.html_url, + prTitle: pr.title, + mergedAt: mergedDate, + isFirstPr, + }); + } + } + } + + return { + handled: true, + prId: prRecordId, + linkedIssueId: matchedIssueId, + registeredAuthor: Boolean(registeredUser), + contributionId, + progressionResult, + }; +} + +/** + * Handles pull_request_review webhook events: + * - Upserts review into `pull_request_reviews` + * - Updates PR state if approved / changes_requested + */ +export async function handlePullRequestReviewWebhook( + payload: WebhookPullRequestReviewEvent, + database?: any +) { + const db = database || (await getDb()); + const review = payload.review; + const now = new Date(); + + // 1. Resolve PR record in database + const matchingPrs = await db + .select() + .from(schema.pullRequests) + .where(eq(schema.pullRequests.url, payload.pull_request.html_url)) + .limit(1); + + let prId: string; + if (matchingPrs.length > 0) { + prId = matchingPrs[0].id; + } else { + // If PR doesn't exist yet, resolve project and insert PR skeleton + const repoFullName = payload.repository.full_name; + const matchingProjects = await db + .select({ id: schema.projects.id }) + .from(schema.projects) + .where(eq(schema.projects.githubRepo, repoFullName)) + .limit(1); + + const projectId = + matchingProjects.length > 0 ? matchingProjects[0].id : `proj_${crypto.randomUUID()}`; + + prId = `pr_${crypto.randomUUID()}`; + await db.insert(schema.pullRequests).values({ + id: prId, + projectId, + githubPrId: payload.pull_request.id, + githubPrNumber: payload.pull_request.number, + title: payload.pull_request.title, + url: payload.pull_request.html_url, + state: "open", + }); + } + + // 2. Map review state + const rawState = review.state.toLowerCase(); + let reviewState: "approved" | "changes_requested" | "commented" | "dismissed" = "commented"; + if (rawState === "approved") { + reviewState = "approved"; + } else if (rawState === "changes_requested") { + reviewState = "changes_requested"; + } else if (rawState === "dismissed") { + reviewState = "dismissed"; + } + + // 3. Upsert into pull_request_reviews + const reviewId = `rev_${crypto.randomUUID()}`; + await db.insert(schema.pullRequestReviews).values({ + id: reviewId, + pullRequestId: prId, + reviewerGithubId: review.user.id, + reviewerUsername: review.user.login, + reviewState, + submittedAt: review.submitted_at ? new Date(review.submitted_at) : now, + githubReviewId: review.id, + htmlUrl: review.html_url, + createdAt: now, + }); + + // 4. Update PR state if approved / changes_requested + if (reviewState === "approved" || reviewState === "changes_requested") { + await db + .update(schema.pullRequests) + .set({ + state: reviewState, + updatedAt: now, + }) + .where(eq(schema.pullRequests.id, prId)); + } + + return { + handled: true, + reviewId, + prId, + reviewState, + }; +}