diff --git a/apps/api/src/controllers/v1/batch-scrape.ts b/apps/api/src/controllers/v1/batch-scrape.ts index ec85b25c3..ecb7cbaed 100644 --- a/apps/api/src/controllers/v1/batch-scrape.ts +++ b/apps/api/src/controllers/v1/batch-scrape.ts @@ -6,7 +6,6 @@ import { batchScrapeRequestSchemaNoURLValidation, url as urlSchema, RequestWithAuth, - ScrapeOptions, BatchScrapeResponse, } from "./types"; import { @@ -19,7 +18,7 @@ import { } from "../../lib/crawl-redis"; import { getJobPriority } from "../../lib/job-priority"; import { addScrapeJobs } from "../../services/queue-jobs"; -import { callWebhook } from "../../services/webhook"; +import { createWebhookSender, WebhookEvent } from "../../services/webhook"; import { logger as _logger } from "../../lib/logger"; import { BLOCKLISTED_URL_MESSAGE } from "../../lib/strings"; import { isUrlBlocked } from "../../scraper/WebScraper/utils/blocklist"; @@ -190,14 +189,12 @@ export async function batchScrapeController( logger.debug("Calling webhook with batch_scrape.started...", { webhook: req.body.webhook, }); - await callWebhook({ + const sender = await createWebhookSender({ teamId: req.auth.team_id, - crawlId: id, - data: null, + jobId: id, webhook: req.body.webhook, - v1: true, - eventType: "batch_scrape.started", }); + await sender?.send(WebhookEvent.BATCH_SCRAPE_STARTED, { success: true }); } const protocol = process.env.ENV === "local" ? req.protocol : "https"; diff --git a/apps/api/src/controllers/v1/extract.ts b/apps/api/src/controllers/v1/extract.ts index fe5e122d8..6a00f466f 100644 --- a/apps/api/src/controllers/v1/extract.ts +++ b/apps/api/src/controllers/v1/extract.ts @@ -6,15 +6,18 @@ import { ExtractResponse, } from "./types"; import { getExtractQueue } from "../../services/queue-service"; -import * as Sentry from "@sentry/node"; import { saveExtract } from "../../lib/extract/extract-redis"; import { getTeamIdSyncB } from "../../lib/extract/team-id-sync"; -import { performExtraction } from "../../lib/extract/extraction-service"; +import { + ExtractResult, + performExtraction, +} from "../../lib/extract/extraction-service"; import { performExtraction_F0 } from "../../lib/extract/fire-0/extraction-service-f0"; import { BLOCKLISTED_URL_MESSAGE } from "../../lib/strings"; import { isUrlBlocked } from "../../scraper/WebScraper/utils/blocklist"; import { logger as _logger } from "../../lib/logger"; import { fromV1ScrapeOptions } from "../v2/types"; +import { createWebhookSender, WebhookEvent } from "../../services/webhook"; export async function oldExtract( req: RequestWithAuth<{}, ExtractResponse, ExtractRequest>, @@ -23,8 +26,17 @@ export async function oldExtract( ) { // Means that are in the non-queue system // TODO: Remove this once all teams have transitioned to the new system + + const sender = await createWebhookSender({ + teamId: req.auth.team_id, + jobId: extractId, + webhook: req.body.webhook, + }); + + sender?.send(WebhookEvent.EXTRACT_STARTED, { success: true }); + try { - let result; + let result: ExtractResult; const model = req.body.agent?.model; if (req.body.agent && model && model.toLowerCase().includes("fire-1")) { result = await performExtraction(extractId, { @@ -42,11 +54,30 @@ export async function oldExtract( }); } + if (sender) { + if (result.success) { + sender.send(WebhookEvent.EXTRACT_COMPLETED, { + success: true, + data: [result], + }); + } else { + sender.send(WebhookEvent.EXTRACT_FAILED, { + success: false, + error: result.error ?? "Unknown error", + }); + } + } + return res.status(200).json(result); } catch (error) { + sender?.send(WebhookEvent.EXTRACT_FAILED, { + success: false, + error: error instanceof Error ? error.message : "Unknown error", + }); + return res.status(500).json({ success: false, - error: "Internal server error", + error: error instanceof Error ? error.message : "Unknown error", }); } } diff --git a/apps/api/src/controllers/v1/types.ts b/apps/api/src/controllers/v1/types.ts index 137e5253b..3b2abdaf0 100644 --- a/apps/api/src/controllers/v1/types.ts +++ b/apps/api/src/controllers/v1/types.ts @@ -13,6 +13,7 @@ import { getURLDepth } from "../../scraper/WebScraper/utils/maxDepthUtils"; import Ajv from "ajv"; import { ErrorCodes } from "../../lib/error"; import { integrationSchema } from "../../utils/integration"; +import { webhookSchema } from "../../services/webhook/schema"; export type Format = | "markdown" @@ -597,6 +598,7 @@ export const extractV1Options = z agent: agentOptionsExtract.optional(), __experimental_showCostTracking: z.boolean().default(false), ignoreInvalidURLs: z.boolean().default(false), + webhook: webhookSchema.optional(), }) .strict(strictMessage) .refine(obj => obj.urls || obj.prompt, { @@ -660,38 +662,6 @@ export const scrapeRequestSchema = baseScrapeOptions export type ScrapeRequest = z.infer; export type ScrapeRequestInput = z.input; -const BLACKLISTED_WEBHOOK_HEADERS = ["x-firecrawl-signature"]; -export const webhookSchema = z.preprocess( - x => { - if (typeof x === "string") { - return { url: x }; - } else { - return x; - } - }, - z - .object({ - url: z.string().url(), - headers: z.record(z.string(), z.string()).default({}), - metadata: z.record(z.string(), z.string()).default({}), - events: z - .array(z.enum(["completed", "failed", "page", "started"])) - .default(["completed", "failed", "page", "started"]), - }) - .strict(strictMessage) - .refine( - obj => { - const blacklistedLower = BLACKLISTED_WEBHOOK_HEADERS.map(h => - h.toLowerCase(), - ); - return !Object.keys(obj.headers).some(key => - blacklistedLower.includes(key.toLowerCase()), - ); - }, - `The following headers are not allowed: ${BLACKLISTED_WEBHOOK_HEADERS.join(", ")}`, - ), -); - export const batchScrapeRequestSchema = baseScrapeOptions .extend({ urls: url.array(), diff --git a/apps/api/src/controllers/v2/batch-scrape.ts b/apps/api/src/controllers/v2/batch-scrape.ts index d211e8024..c00d9a0e7 100644 --- a/apps/api/src/controllers/v2/batch-scrape.ts +++ b/apps/api/src/controllers/v2/batch-scrape.ts @@ -19,7 +19,7 @@ import { } from "../../lib/crawl-redis"; import { getJobPriority } from "../../lib/job-priority"; import { addScrapeJobs } from "../../services/queue-jobs"; -import { callWebhook } from "../../services/webhook"; +import { createWebhookSender, WebhookEvent } from "../../services/webhook"; import { logger as _logger } from "../../lib/logger"; import { BLOCKLISTED_URL_MESSAGE } from "../../lib/strings"; import { isUrlBlocked } from "../../scraper/WebScraper/utils/blocklist"; @@ -186,14 +186,12 @@ export async function batchScrapeController( logger.debug("Calling webhook with batch_scrape.started...", { webhook: req.body.webhook, }); - await callWebhook({ + const sender = await createWebhookSender({ teamId: req.auth.team_id, - crawlId: id, - data: null, + jobId: id, webhook: req.body.webhook, - v1: true, - eventType: "batch_scrape.started", }); + await sender?.send(WebhookEvent.BATCH_SCRAPE_STARTED, { success: true }); } const protocol = process.env.ENV === "local" ? req.protocol : "https"; diff --git a/apps/api/src/controllers/v2/types.ts b/apps/api/src/controllers/v2/types.ts index 1d624daf1..89ad8cfbd 100644 --- a/apps/api/src/controllers/v2/types.ts +++ b/apps/api/src/controllers/v2/types.ts @@ -18,6 +18,7 @@ import type { InternalOptions } from "../../scraper/scrapeURL"; import { ErrorCodes } from "../../lib/error"; import Ajv from "ajv"; import { integrationSchema } from "../../utils/integration"; +import { webhookSchema } from "../../services/webhook/schema"; export type Format = | "markdown" @@ -519,6 +520,7 @@ export const extractOptions = z .optional(), __experimental_showCostTracking: z.boolean().default(false), ignoreInvalidURLs: z.boolean().default(true), + webhook: webhookSchema.optional(), }) .strict(strictMessage) .refine(obj => obj.urls || obj.prompt, { @@ -558,38 +560,6 @@ export const scrapeRequestSchema = baseScrapeOptions export type ScrapeRequest = z.infer; export type ScrapeRequestInput = z.input; -const BLACKLISTED_WEBHOOK_HEADERS = ["x-firecrawl-signature"]; -export const webhookSchema = z.preprocess( - x => { - if (typeof x === "string") { - return { url: x }; - } else { - return x; - } - }, - z - .object({ - url: z.string().url(), - headers: z.record(z.string(), z.string()).default({}), - metadata: z.record(z.string(), z.string()).default({}), - events: z - .array(z.enum(["completed", "failed", "page", "started"])) - .default(["completed", "failed", "page", "started"]), - }) - .strict(strictMessage) - .refine( - obj => { - const blacklistedLower = BLACKLISTED_WEBHOOK_HEADERS.map(h => - h.toLowerCase(), - ); - return !Object.keys(obj.headers).some(key => - blacklistedLower.includes(key.toLowerCase()), - ); - }, - `The following headers are not allowed: ${BLACKLISTED_WEBHOOK_HEADERS.join(", ")}`, - ), -); - export const batchScrapeRequestSchema = baseScrapeOptions .extend({ urls: url.array(), diff --git a/apps/api/src/services/queue-worker.ts b/apps/api/src/services/queue-worker.ts index 09f7352dc..6b3fae215 100644 --- a/apps/api/src/services/queue-worker.ts +++ b/apps/api/src/services/queue-worker.ts @@ -33,6 +33,7 @@ import { performDeepResearch } from "../lib/deep-research/deep-research-service" import { performGenerateLlmsTxt } from "../lib/generate-llmstxt/generate-llmstxt-service"; import { updateGeneratedLlmsTxt } from "../lib/generate-llmstxt/generate-llmstxt-redis"; import { performExtraction_F0 } from "../lib/extract/fire-0/extraction-service-f0"; +import { createWebhookSender, WebhookEvent } from "./webhook"; import Express from "express"; import http from "http"; import https from "https"; @@ -118,7 +119,20 @@ const processExtractJobInternal = async ( await job.extendLock(token, jobLockExtensionTime); }, jobLockExtendInterval); + const sender = await createWebhookSender({ + teamId: job.data.teamId, + jobId: job.data.extractId, + webhook: job.data.request.webhook, + v0: true, + }); + try { + if (sender) { + sender.send(WebhookEvent.EXTRACT_STARTED, { + success: true, + }); + } + let result: ExtractResult | null = null; const model = job.data.request.agent?.model; @@ -150,6 +164,14 @@ const processExtractJobInternal = async ( if (result && result.success) { // Move job to completed state in Redis await job.moveToCompleted(result, token, false); + + if (sender) { + sender.send(WebhookEvent.EXTRACT_COMPLETED, { + success: true, + data: [result], + }); + } + return result; } else { // throw new Error(result.error || "Unknown error during extraction"); @@ -163,6 +185,16 @@ const processExtractJobInternal = async ( job.data.extractId, }); + if (sender) { + sender.send(WebhookEvent.EXTRACT_FAILED, { + success: false, + error: + result?.error ?? + "Unknown error, please contact help@firecrawl.com. Extract id: " + + job.data.extractId, + }); + } + return result; } } catch (error) { @@ -189,6 +221,17 @@ const processExtractJobInternal = async ( "Unknown error, please contact help@firecrawl.com. Extract id: " + job.data.extractId, }); + + if (sender) { + sender.send(WebhookEvent.EXTRACT_FAILED, { + success: false, + error: + (error as any)?.message ?? + "Unknown error, please contact help@firecrawl.com. Extract id: " + + job.data.extractId, + }); + } + return { success: false, error: diff --git a/apps/api/src/services/webhook.ts b/apps/api/src/services/webhook.ts deleted file mode 100644 index 5b64a19c4..000000000 --- a/apps/api/src/services/webhook.ts +++ /dev/null @@ -1,330 +0,0 @@ -import undici from "undici"; -import { logger as _logger, logger } from "../lib/logger"; -import { supabase_rr_service, supabase_service } from "./supabase"; -import { WebhookEventType } from "../types"; -import { configDotenv } from "dotenv"; -import { z } from "zod"; -import { webhookSchema } from "../controllers/v1/types"; -import { redisEvictConnection } from "./redis"; -import { createHmac } from "crypto"; -import { - getSecureDispatcher, - isIPPrivate, -} from "../scraper/scrapeURL/engines/utils/safeFetch"; -configDotenv(); - -const WEBHOOK_INSERT_QUEUE_KEY = "webhook-insert-queue"; -const WEBHOOK_INSERT_BATCH_SIZE = 1000; - -interface WebhookLogData { - success: boolean; - error?: string; - teamId: string; - crawlId: string; - scrapeId?: string; - url: string; - statusCode?: number; - event: WebhookEventType; -} - -interface CallWebhookParams { - teamId: string; - crawlId: string; - scrapeId?: string; - webhook?: z.infer; - v1: boolean; - data: any | null; - eventType: WebhookEventType; - awaitWebhook?: boolean; -} - -async function addWebhookInsertJob(data: any) { - await redisEvictConnection.rpush( - WEBHOOK_INSERT_QUEUE_KEY, - JSON.stringify(data), - ); -} - -export async function getWebhookInsertQueueLength(): Promise { - return (await redisEvictConnection.llen(WEBHOOK_INSERT_QUEUE_KEY)) ?? 0; -} - -async function getWebhookInsertJobs(): Promise { - const jobs = - (await redisEvictConnection.lpop( - WEBHOOK_INSERT_QUEUE_KEY, - WEBHOOK_INSERT_BATCH_SIZE, - )) ?? []; - return jobs.map(x => JSON.parse(x)); -} - -export async function processWebhookInsertJobs() { - const jobs = await getWebhookInsertJobs(); - if (jobs.length === 0) { - return; - } - logger.info(`Webhook inserter found jobs to insert`, { - jobCount: jobs.length, - }); - try { - await supabase_service.from("webhook_logs").insert(jobs); - logger.info(`Webhook inserter inserted jobs`, { jobCount: jobs.length }); - } catch (error) { - logger.error(`Webhook inserter failed to insert jobs`, { - error, - jobCount: jobs.length, - }); - } -} - -async function logWebhook(data: WebhookLogData) { - try { - await addWebhookInsertJob({ - success: data.success, - error: data.error ?? null, - team_id: data.teamId, - crawl_id: data.crawlId, - scrape_id: data.scrapeId ?? null, - url: data.url, - status_code: data.statusCode ?? null, - event: data.event, - }); - } catch (error) { - _logger.error("Error logging webhook", { - error, - crawlId: data.crawlId, - scrapeId: data.scrapeId, - teamId: data.teamId, - team_id: data.teamId, - module: "webhook", - method: "logWebhook", - }); - } -} - -function generateHmacSignature(payload: string, secret: string): string { - const hmac = createHmac("sha256", secret); - hmac.update(payload); - return hmac.digest("hex"); -} - -export const callWebhook = async ({ - teamId, - crawlId, - scrapeId, - data, - webhook, - v1, - eventType, - awaitWebhook = false, -}: CallWebhookParams) => { - const logger = _logger.child({ - module: "webhook", - method: "callWebhook", - teamId, - team_id: teamId, - crawlId, - scrapeId, - eventType, - awaitWebhook, - webhook, - isV1: v1, - }); - - if (webhook) { - let subType = eventType.split(".")[1]; - if (!webhook.events.includes(subType as any)) { - logger.debug("Webhook event type not in specified events", { - subType, - webhook, - }); - return false; - } - } - - try { - const selfHostedUrl = process.env.SELF_HOSTED_WEBHOOK_URL?.replace( - "{{JOB_ID}}", - crawlId, - ); - const useDbAuthentication = process.env.USE_DB_AUTHENTICATION === "true"; - let webhookUrl = - webhook ?? - (selfHostedUrl ? webhookSchema.parse({ url: selfHostedUrl }) : undefined); - - // TODO: do these queries upstream so we don't have to query the DB multiple times per crawl - - // Only fetch the webhook URL from the database if the self-hosted webhook URL and specified webhook are not set - // and the USE_DB_AUTHENTICATION environment variable is set to true - if (!webhookUrl && useDbAuthentication) { - const { data: webhooksData, error } = await supabase_rr_service - .from("webhooks") - .select("url") - .eq("team_id", teamId) - .limit(1); - if (error) { - logger.error(`Error fetching webhook URL for team`, { - error, - }); - return null; - } - - if (!webhooksData || webhooksData.length === 0) { - return null; - } - - webhookUrl = webhooksData[0].url; - } - - let hmacSecret: string | undefined = - process.env.SELF_HOSTED_WEBHOOK_HMAC_SECRET; - if (useDbAuthentication) { - const { data: teamData, error: teamError } = await supabase_rr_service - .from("teams") - .select("hmac_secret") - .eq("id", teamId) - .limit(1) - .single(); - - if (teamError) { - logger.error(`Error fetching team HMAC secret`, { - error: teamError, - }); - } - - if (teamData?.hmac_secret) hmacSecret = teamData.hmac_secret; - } - - logger.debug("Calling webhook...", { - webhookUrl, - }); - - if (!webhookUrl) { - return null; - } - - // check if the webhook URL is a private IP address *before* making the request - // the dispatcher also performs a check once connected, however this prevents unnecessary connections - const webhookHost = new URL(webhookUrl.url).hostname; - if (isIPPrivate(webhookHost)) { - logger.warn("Aborting webhook call to private IP address", { - url: webhookUrl.url, - }); - return null; - } - - let dataToSend: any[] = []; - if ( - data && - data.result && - data.result.links && - data.result.links.length !== 0 - ) { - for (let i = 0; i < data.result.links.length; i++) { - if (v1) { - dataToSend.push(data.result.links[i].content); - } else { - dataToSend.push({ - content: data.result.links[i].content.content, - markdown: data.result.links[i].content.markdown, - metadata: data.result.links[i].content.metadata, - }); - } - } - } - - const payload = { - success: !v1 - ? data.success - : eventType === "crawl.page" - ? data.success - : true, - type: eventType, - [v1 ? "id" : "jobId"]: crawlId, - data: dataToSend, - error: !v1 - ? data?.error || undefined - : eventType === "crawl.page" - ? data?.error || undefined - : undefined, - metadata: webhookUrl.metadata || undefined, - }; - - const payloadString = JSON.stringify(payload); - - const headers: Record = { - "Content-Type": "application/json", - ...webhookUrl.headers, - }; - - if (hmacSecret) { - const signature = generateHmacSignature(payloadString, hmacSecret); - headers["X-Firecrawl-Signature"] = `sha256=${signature}`; - } - - const executeWebhookRequest = async () => { - try { - const res = await undici.fetch(webhookUrl.url, { - method: "POST", - headers, - body: payloadString, - dispatcher: getSecureDispatcher(), - signal: AbortSignal.timeout(v1 ? 10000 : 30000), // 10 seconds timeout (v1) - }); - - if (!res.ok) { - throw { status: res.status }; - } - - logWebhook({ - success: res.status >= 200 && res.status < 300, - teamId, - crawlId, - scrapeId, - url: webhookUrl.url, - event: eventType, - statusCode: res.status, - }); - - return res; - } catch (error) { - logger.error(`Failed to send webhook`, { - error, - }); - - logWebhook({ - success: false, - teamId, - crawlId, - scrapeId, - url: webhookUrl.url, - event: eventType, - error: - error instanceof Error - ? error.message - : typeof error === "string" - ? error - : undefined, - statusCode: - typeof (error as any)?.status === "number" - ? (error as any).status - : undefined, - }); - - throw error; - } - }; - - if (awaitWebhook) { - await executeWebhookRequest(); - } else { - executeWebhookRequest().catch(() => { - // already logged in executeWebhookRequest - }); - } - } catch (error) { - logger.warn(`Error sending webhook`, { - error, - }); - } -}; diff --git a/apps/api/src/services/webhook/config.ts b/apps/api/src/services/webhook/config.ts new file mode 100644 index 000000000..36f1dc784 --- /dev/null +++ b/apps/api/src/services/webhook/config.ts @@ -0,0 +1,78 @@ +import { logger as _logger } from "../../lib/logger"; +import { supabase_rr_service } from "../supabase"; +import { WebhookConfig } from "./types"; + +export async function getWebhookConfig( + teamId: string, + jobId: string, + webhook?: WebhookConfig, +): Promise<{ config: WebhookConfig; secret?: string } | null> { + // priority: + // - webhook + // - self-hosted environment variable + // - db webhook (if enabled) + if (webhook) { + return { config: webhook, secret: await getHmacSecret(teamId) }; + } + + const selfHostedUrl = process.env.SELF_HOSTED_WEBHOOK_URL?.replace( + "{{JOB_ID}}", + jobId, + ); + if (selfHostedUrl) { + return { + config: { + url: selfHostedUrl, + headers: {}, + metadata: {}, + events: ["completed", "failed", "page", "started"], + }, + secret: process.env.SELF_HOSTED_WEBHOOK_HMAC_SECRET, + }; + } + + if (process.env.USE_DB_AUTHENTICATION === "true") { + const dbConfig = await fetchWebhookFromDb(teamId); + if (dbConfig) { + return { config: dbConfig, secret: await getHmacSecret(teamId) }; + } + } + + return null; +} + +async function fetchWebhookFromDb( + teamId: string, +): Promise { + try { + const { data, error } = await supabase_rr_service + .from("webhooks") + .select("url, headers, metadata, events") + .eq("team_id", teamId) + .limit(1) + .single(); + + return error || !data ? null : data; + } catch { + return null; + } +} + +async function getHmacSecret(teamId: string): Promise { + if (process.env.USE_DB_AUTHENTICATION !== "true") { + return process.env.SELF_HOSTED_WEBHOOK_HMAC_SECRET; + } + + try { + const { data, error } = await supabase_rr_service + .from("teams") + .select("hmac_secret") + .eq("id", teamId) + .limit(1) + .single(); + + return error ? undefined : data?.hmac_secret; + } catch { + return undefined; + } +} diff --git a/apps/api/src/services/webhook/delivery.ts b/apps/api/src/services/webhook/delivery.ts new file mode 100644 index 000000000..4c9448d9d --- /dev/null +++ b/apps/api/src/services/webhook/delivery.ts @@ -0,0 +1,204 @@ +import undici from "undici"; +import { createHmac } from "crypto"; +import { logger as _logger, logger } from "../../lib/logger"; +import { + getSecureDispatcher, + isIPPrivate, +} from "../../scraper/scrapeURL/engines/utils/safeFetch"; +import { WebhookConfig, WebhookEvent, WebhookEventDataMap } from "./types"; +import { redisEvictConnection } from "../redis"; +import { supabase_service } from "../supabase"; + +const WEBHOOK_INSERT_QUEUE_KEY = "webhook-insert-queue"; +const WEBHOOK_INSERT_BATCH_SIZE = 1000; + +export class WebhookSender { + private config: WebhookConfig; + private secret?: string; + private context: { teamId: string; jobId: string; v0: boolean }; + private logger: any; + + constructor( + config: WebhookConfig, + secret: string | undefined, + context: { teamId: string; jobId: string; v0: boolean }, + ) { + this.config = config; + this.secret = secret; + this.context = context; + this.logger = _logger.child({ + module: "webhook-sender", + teamId: context.teamId, + jobId: context.jobId, + isV0: context.v0, + }); + } + + async send( + event: T, + data: WebhookEventDataMap[T], + ): Promise { + if (!this.shouldSendEvent(event)) return; + + const payload = { + success: data.success, + type: event, + [this.context.v0 ? "jobId" : "id"]: this.context.jobId, + data: "data" in data ? data.data : [], + error: "error" in data ? data.error : undefined, + metadata: this.config.metadata || undefined, + }; + + const delivery = this.deliver(payload, data.scrapeId); + + if (data.awaitWebhook) { + await delivery; + } else { + delivery.catch(() => {}); + } + } + + private shouldSendEvent(event: WebhookEvent): boolean { + if (!this.config.events?.length) return true; + const subType = event.split(".")[1]; + return this.config.events.includes(subType as any); + } + + private async deliver(payload: any, scrapeId?: string): Promise { + const webhookHost = new URL(this.config.url).hostname; + if (isIPPrivate(webhookHost)) { + this.logger.warn("Aborting webhook call to private IP address", { + webhookUrl: this.config.url, + }); + return; + } + + const payloadString = JSON.stringify(payload); + const headers: Record = { + "Content-Type": "application/json", + ...this.config.headers, + }; + + if (this.secret) { + const hmac = createHmac("sha256", this.secret); + hmac.update(payloadString); + headers["X-Firecrawl-Signature"] = `sha256=${hmac.digest("hex")}`; + } + + try { + const res = await undici.fetch(this.config.url, { + method: "POST", + headers, + body: payloadString, + dispatcher: getSecureDispatcher(), + signal: AbortSignal.timeout(this.context.v0 ? 30000 : 10000), + }); + + await logWebhook({ + success: res.status >= 200 && res.status < 300, + teamId: this.context.teamId, + crawlId: this.context.jobId, // this is legacy naming, we should rename it to jobId at some point + scrapeId, + url: this.config.url, + event: payload.type, + statusCode: res.status, + }); + + if (!res.ok) { + throw new Error(`Unexpected response status: ${res.status}`); + } + } catch (error) { + this.logger.error("Failed to send webhook", { + error, + webhookUrl: this.config.url, + }); + + await logWebhook({ + success: false, + teamId: this.context.teamId, + crawlId: this.context.jobId, // same as above + scrapeId, + url: this.config.url, + event: payload.type, + error: + error instanceof Error + ? error.message + : typeof error === "string" + ? error + : undefined, + statusCode: + typeof (error as any)?.status === "number" + ? (error as any).status + : undefined, + }); + + throw error; + } + } + + get metadata(): Record { + return this.config.metadata || {}; + } +} + +export async function getWebhookInsertQueueLength(): Promise { + return (await redisEvictConnection.llen(WEBHOOK_INSERT_QUEUE_KEY)) ?? 0; +} + +export async function processWebhookInsertJobs() { + const jobs = + (await redisEvictConnection.lpop( + WEBHOOK_INSERT_QUEUE_KEY, + WEBHOOK_INSERT_BATCH_SIZE, + )) ?? []; + if (jobs.length === 0) return; + + const parsedJobs = jobs.map(x => JSON.parse(x)); + _logger.info("Webhook inserter found jobs to insert", { + jobCount: parsedJobs.length, + }); + + try { + await supabase_service.from("webhook_logs").insert(parsedJobs); + _logger.info("Webhook inserter inserted jobs", { + jobCount: parsedJobs.length, + }); + } catch (error) { + _logger.error("Webhook inserter failed to insert jobs", { + error, + jobCount: parsedJobs.length, + }); + } +} + +export async function logWebhook(data: { + success: boolean; + error?: string; + teamId: string; + crawlId: string; + scrapeId?: string; + url: string; + statusCode?: number; + event: WebhookEvent; +}): Promise { + try { + await redisEvictConnection.rpush( + WEBHOOK_INSERT_QUEUE_KEY, + JSON.stringify({ + success: data.success, + error: data.error ?? null, + team_id: data.teamId, + crawl_id: data.crawlId, + scrape_id: data.scrapeId ?? null, + url: data.url, + status_code: data.statusCode ?? null, + event: data.event, + }), + ); + } catch (error) { + logger.error("Error logging webhook", { + error, + teamId: data.teamId, + }); + } +} diff --git a/apps/api/src/services/webhook/index.ts b/apps/api/src/services/webhook/index.ts new file mode 100644 index 000000000..e8d20706f --- /dev/null +++ b/apps/api/src/services/webhook/index.ts @@ -0,0 +1,32 @@ +import { logger as _logger } from "../../lib/logger"; +import { getWebhookConfig } from "./config"; +import { WebhookConfig } from "./types"; +import { WebhookSender } from "./delivery"; + +export async function createWebhookSender(params: { + teamId: string; + jobId: string; + webhook?: WebhookConfig; + v0?: boolean; +}): Promise { + const config = await getWebhookConfig( + params.teamId, + params.jobId, + params.webhook, + ); + if (!config) { + return null; + } + + return new WebhookSender(config.config, config.secret, { + teamId: params.teamId, + jobId: params.jobId, + v0: params.v0 || false, + }); +} + +export { + getWebhookInsertQueueLength, + processWebhookInsertJobs, +} from "./delivery"; +export * from "./types"; diff --git a/apps/api/src/services/webhook/schema.ts b/apps/api/src/services/webhook/schema.ts new file mode 100644 index 000000000..f4c1627f4 --- /dev/null +++ b/apps/api/src/services/webhook/schema.ts @@ -0,0 +1,30 @@ +import { z } from "zod"; + +const BLACKLISTED_WEBHOOK_HEADERS = ["x-firecrawl-signature"]; + +export const webhookSchema = z.preprocess( + x => (typeof x === "string" ? { url: x } : x), + z + .object({ + url: z.string().url(), + headers: z.record(z.string(), z.string()).default({}), + metadata: z.record(z.string(), z.string()).default({}), + events: z + .array(z.enum(["completed", "failed", "page", "started"])) + .default(["completed", "failed", "page", "started"]), + }) + .strict( + "Unrecognized key in webhook object. Review the API documentation for webhook configuration changes.", + ) + .refine( + obj => { + const blacklistedLower = BLACKLISTED_WEBHOOK_HEADERS.map(h => + h.toLowerCase(), + ); + return !Object.keys(obj.headers).some(key => + blacklistedLower.includes(key.toLowerCase()), + ); + }, + `The following headers are not allowed: ${BLACKLISTED_WEBHOOK_HEADERS.join(", ")}`, + ), +); diff --git a/apps/api/src/services/webhook/types.ts b/apps/api/src/services/webhook/types.ts new file mode 100644 index 000000000..39e33689b --- /dev/null +++ b/apps/api/src/services/webhook/types.ts @@ -0,0 +1,98 @@ +import { z } from "zod"; +import { webhookSchema } from "./schema"; +import { ExtractResult } from "../../lib/extract/extraction-service"; + +export enum WebhookEvent { + CRAWL_STARTED = "crawl.started", + CRAWL_PAGE = "crawl.page", + CRAWL_COMPLETED = "crawl.completed", + BATCH_SCRAPE_STARTED = "batch_scrape.started", + BATCH_SCRAPE_PAGE = "batch_scrape.page", + BATCH_SCRAPE_COMPLETED = "batch_scrape.completed", + EXTRACT_STARTED = "extract.started", + EXTRACT_COMPLETED = "extract.completed", + EXTRACT_FAILED = "extract.failed", +} + +export type WebhookEventDataMap = { + [WebhookEvent.CRAWL_STARTED]: CrawlStartedData; + [WebhookEvent.CRAWL_PAGE]: CrawlPageData; + [WebhookEvent.CRAWL_COMPLETED]: CrawlCompletedData; + [WebhookEvent.BATCH_SCRAPE_STARTED]: BatchScrapeStartedData; + [WebhookEvent.BATCH_SCRAPE_PAGE]: BatchScrapePageData; + [WebhookEvent.BATCH_SCRAPE_COMPLETED]: BatchScrapeCompletedData; + [WebhookEvent.EXTRACT_STARTED]: ExtractStartedData; + [WebhookEvent.EXTRACT_COMPLETED]: ExtractCompletedData; + [WebhookEvent.EXTRACT_FAILED]: ExtractFailedData; +}; + +export type WebhookConfig = z.infer; + +export interface WebhookDocument { + content?: string; + markdown: string; + metadata: Record; +} + +export interface WebhookDocumentLink { + content: WebhookDocument; + source: string; +} + +interface BaseWebhookData { + success: boolean; + scrapeId?: string; + awaitWebhook?: boolean; +} + +// crawl +export interface CrawlStartedData extends BaseWebhookData { + success: true; +} + +export interface CrawlPageData extends BaseWebhookData { + success: boolean; + data: WebhookDocument[] | WebhookDocumentLink[]; // links or documents (v0 compatible) + error?: string; +} + +export interface CrawlCompletedData extends BaseWebhookData { + success: true; + data: WebhookDocument[] | WebhookDocumentLink[]; // empty array or links (v0 compatible) +} + +export interface CrawlFailedData extends BaseWebhookData { + success: false; + error: string; +} + +// batch scrape +export interface BatchScrapeStartedData extends BaseWebhookData { + success: true; +} + +export interface BatchScrapePageData extends BaseWebhookData { + success: boolean; + data: WebhookDocument[]; + error?: string; // more v0 tomfoolery +} + +export interface BatchScrapeCompletedData extends BaseWebhookData { + success: true; + data: WebhookDocumentLink[]; +} + +// extract +export interface ExtractStartedData extends BaseWebhookData { + success: true; +} + +export interface ExtractCompletedData extends BaseWebhookData { + success: true; + data: ExtractResult[]; +} + +export interface ExtractFailedData extends BaseWebhookData { + success: false; + error: string; +} diff --git a/apps/api/src/services/worker/crawl-logic.ts b/apps/api/src/services/worker/crawl-logic.ts index 974b431e3..677dea287 100644 --- a/apps/api/src/services/worker/crawl-logic.ts +++ b/apps/api/src/services/worker/crawl-logic.ts @@ -19,7 +19,7 @@ import { v4 as uuidv4 } from "uuid"; import { addScrapeJobs } from "../queue-jobs"; import { getJobs } from "../../controllers/v1/crawl-status"; import { logJob } from "../logging/log_job"; -import { callWebhook } from "../webhook"; +import { createWebhookSender, WebhookEvent } from "../webhook"; import { hasFormatOfType } from "../../lib/format-utils"; export async function finishCrawlIfNeeded( @@ -231,17 +231,33 @@ export async function finishCrawlIfNeeded( // v0 web hooks, call when done with all the data if (!job.data.v1) { - callWebhook({ + const sender = await createWebhookSender({ teamId: job.data.team_id, - crawlId: job.data.crawl_id, - data, - webhook: job.data.webhook, - v1: job.data.v1, - eventType: - job.data.crawlerOptions !== null - ? "crawl.completed" - : "batch_scrape.completed", + jobId: job.data.crawl_id, + webhook: job.data.webhook as any, + v0: true, }); + if (sender) { + const documents = fullDocs.map((doc: any) => ({ + content: { + content: doc?.content ?? doc?.rawHtml ?? doc?.markdown ?? "", + markdown: doc?.markdown, + metadata: doc?.metadata ?? {}, + }, + source: doc?.metadata?.sourceURL ?? doc?.url ?? "", + })); + if (job.data.crawlerOptions !== null) { + sender.send(WebhookEvent.CRAWL_COMPLETED, { + success: true, + data: documents, + }); + } else { + sender.send(WebhookEvent.BATCH_SCRAPE_COMPLETED, { + success: true, + data: documents, + }); + } + } } } else { const num_docs = await getDoneJobsOrderedLength(job.data.crawl_id); @@ -292,17 +308,24 @@ export async function finishCrawlIfNeeded( // v1 web hooks, call when done with no data, but with event completed if (job.data.v1 && job.data.webhook) { - callWebhook({ + const sender = await createWebhookSender({ teamId: job.data.team_id, - crawlId: job.data.crawl_id, - data: [], - webhook: job.data.webhook, - v1: job.data.v1, - eventType: - job.data.crawlerOptions !== null - ? "crawl.completed" - : "batch_scrape.completed", + jobId: job.data.crawl_id, + webhook: job.data.webhook as any, }); + if (sender) { + if (job.data.crawlerOptions !== null) { + sender.send(WebhookEvent.CRAWL_COMPLETED, { + success: true, + data: [], + }); + } else { + sender.send(WebhookEvent.BATCH_SCRAPE_COMPLETED, { + success: true, + data: [], + }); + } + } } } } diff --git a/apps/api/src/services/worker/scrape-worker.ts b/apps/api/src/services/worker/scrape-worker.ts index 1c79ccd4a..69fd6aa7b 100644 --- a/apps/api/src/services/worker/scrape-worker.ts +++ b/apps/api/src/services/worker/scrape-worker.ts @@ -31,7 +31,7 @@ import { getJobPriority } from "../../lib/job-priority"; import { Document, scrapeOptions, TeamFlags } from "../../controllers/v2/types"; import { hasFormatOfType } from "../../lib/format-utils"; import { getACUCTeam } from "../../controllers/auth"; -import { callWebhook } from "../webhook"; +import { createWebhookSender, WebhookEvent } from "../webhook"; import { CustomError } from "../../lib/custom-error"; import { startWebScraperPipeline } from "../../main/runWebScraper"; import { CostTracking } from "../../lib/cost-tracking"; @@ -453,18 +453,31 @@ async function processJob(job: Job & { id: string }) { logger.debug("Calling webhook with success...", { webhook: job.data.webhook, }); - callWebhook({ + const sender = await createWebhookSender({ teamId: job.data.team_id, - crawlId: job.data.crawl_id, - scrapeId: job.id, - data, - webhook: job.data.webhook, - v1: job.data.v1, - eventType: - job.data.crawlerOptions !== null - ? "crawl.page" - : "batch_scrape.page", + jobId: job.data.crawl_id, + webhook: job.data.webhook as any, }); + if (sender) { + const documents = Array.isArray(data?.result?.links) + ? data.result.links.map((x: any) => ({ + content: x?.content?.content, + markdown: x?.content?.markdown, + metadata: x?.content?.metadata, + })) + : []; + if (job.data.crawlerOptions !== null) { + sender.send(WebhookEvent.CRAWL_PAGE, { + success: true, + data: documents, + }); + } else { + sender.send(WebhookEvent.BATCH_SCRAPE_PAGE, { + success: true, + data: documents, + }); + } + } } logger.debug("Declaring job as done..."); @@ -580,16 +593,33 @@ async function processJob(job: Job & { id: string }) { }; if (!job.data.v1 && (job.data.mode === "crawl" || job.data.crawl_id)) { - callWebhook({ + const sender = await createWebhookSender({ teamId: job.data.team_id, - crawlId: job.data.crawl_id ?? (job.id as string), - scrapeId: job.id, - data, - webhook: job.data.webhook, - v1: job.data.v1, - eventType: - job.data.crawlerOptions !== null ? "crawl.page" : "batch_scrape.page", + jobId: (job.data.crawl_id ?? (job.id as string)) as string, + webhook: job.data.webhook as any, + v0: true, }); + if (sender) { + const errorMessage = + data?.error instanceof Error + ? data.error.message + : typeof data?.error === "string" + ? data.error + : "Unknown error"; + if (job.data.crawlerOptions !== null) { + sender.send(WebhookEvent.CRAWL_PAGE, { + success: false, + error: errorMessage, + data: [], + }); + } else { + sender.send(WebhookEvent.BATCH_SCRAPE_PAGE, { + success: false, + error: errorMessage, + data: [], + }); + } + } } const end = Date.now(); @@ -717,14 +747,15 @@ async function processKickoffJob(job: Job & { id: string }) { logger.debug("Calling webhook with crawl.started...", { webhook: job.data.webhook, }); - callWebhook({ + const sender = await createWebhookSender({ teamId: job.data.team_id, - crawlId: job.data.crawl_id, - data: null, - webhook: job.data.webhook, - v1: job.data.v1, - eventType: "crawl.started", + jobId: job.data.crawl_id, + webhook: job.data.webhook as any, + v0: Boolean(!job.data.v1), }); + if (sender) { + sender.send(WebhookEvent.CRAWL_STARTED, { success: true }); + } } const sitemap = sc.crawlerOptions.ignoreSitemap diff --git a/apps/api/src/types.ts b/apps/api/src/types.ts index 834e81dd7..b5f99a7e9 100644 --- a/apps/api/src/types.ts +++ b/apps/api/src/types.ts @@ -3,13 +3,13 @@ import { BaseScrapeOptions, ScrapeOptions, Document as V2Document, - webhookSchema, TeamFlags, } from "./controllers/v2/types"; import { AuthCreditUsageChunk } from "./controllers/v1/types"; import { ExtractorOptions, Document } from "./lib/entities"; import { InternalOptions } from "./scraper/scrapeURL"; import type { CostTracking } from "./lib/cost-tracking"; +import { webhookSchema } from "./services/webhook/schema"; type Mode = "crawl" | "single_urls" | "sitemap" | "kickoff"; @@ -201,12 +201,3 @@ export type ScrapeLog = { ipv4_support?: boolean | null; ipv6_support?: boolean | null; }; - -export type WebhookEventType = - | "crawl.page" - | "batch_scrape.page" - | "crawl.started" - | "batch_scrape.started" - | "crawl.completed" - | "batch_scrape.completed" - | "crawl.failed";