From 24f06ed181d8b72840107d2baef8fc139a89d811 Mon Sep 17 00:00:00 2001 From: Nicolas Date: Sat, 27 Sep 2025 12:02:21 +0300 Subject: [PATCH] (feat/big-query) Big Query (#2217) * Nick: big query * Update log_job.ts --- apps/api/package.json | 1 + apps/api/pnpm-lock.yaml | 193 ++++++++++++ apps/api/src/lib/bigquery-jobs.ts | 370 +++++++++++++++++++++++ apps/api/src/lib/job-transform.ts | 131 ++++++++ apps/api/src/services/logging/log_job.ts | 118 ++------ 5 files changed, 720 insertions(+), 93 deletions(-) create mode 100644 apps/api/src/lib/bigquery-jobs.ts create mode 100644 apps/api/src/lib/job-transform.ts diff --git a/apps/api/package.json b/apps/api/package.json index a384f6757..8d88296a1 100644 --- a/apps/api/package.json +++ b/apps/api/package.json @@ -73,6 +73,7 @@ "@bull-board/express": "^6.11.2", "@coinbase/x402": "^0.6.4", "@dqbd/tiktoken": "^1.0.22", + "@google-cloud/bigquery": "^8.1.1", "@google-cloud/storage": "^7.16.0", "@mendable/firecrawl-rs": "workspace:*", "@openrouter/ai-sdk-provider": "^0.4.5", diff --git a/apps/api/pnpm-lock.yaml b/apps/api/pnpm-lock.yaml index 4ec30a3d8..c4c67b301 100644 --- a/apps/api/pnpm-lock.yaml +++ b/apps/api/pnpm-lock.yaml @@ -48,6 +48,9 @@ importers: '@dqbd/tiktoken': specifier: ^1.0.22 version: 1.0.22 + '@google-cloud/bigquery': + specifier: ^8.1.1 + version: 8.1.1 '@google-cloud/storage': specifier: ^7.16.0 version: 7.16.0(encoding@0.1.13) @@ -961,10 +964,26 @@ packages: resolution: {integrity: sha512-zQ0IqbdX8FZ9aw11vP+dZkKDkS+kgIvQPHnSAXzP9pLu+Rfu3D3XEeLbicvoXJTYnhZiPmsZUxgdzXwNKxRPbA==} engines: {node: '>=14'} + '@google-cloud/bigquery@8.1.1': + resolution: {integrity: sha512-2GHlohfA/VJffTvibMazMsZi6jPRx8MmaMberyDTL8rnhVs/frKSXVVRtLU83uSAy2j/5SD4mOs4jMQgJPON2g==} + engines: {node: '>=18'} + + '@google-cloud/common@6.0.0': + resolution: {integrity: sha512-IXh04DlkLMxWgYLIUYuHHKXKOUwPDzDgke1ykkkJPe48cGIS9kkL2U/o0pm4ankHLlvzLF/ma1eO86n/bkumIA==} + engines: {node: '>=18'} + '@google-cloud/paginator@5.0.2': resolution: {integrity: sha512-DJS3s0OVH4zFDB1PzjxAsHqJT6sKVbRwwML0ZBP9PbU7Yebtu/7SWMRzvO2J3nUi9pRNITCfu4LJeooM2w4pjg==} engines: {node: '>=14.0.0'} + '@google-cloud/paginator@6.0.0': + resolution: {integrity: sha512-g5nmMnzC+94kBxOKkLGpK1ikvolTFCC3s2qtE4F+1EuArcJ7HHC23RDQVt3Ra3CqpUYZ+oXNKZ8n5Cn5yug8DA==} + engines: {node: '>=18'} + + '@google-cloud/precise-date@5.0.0': + resolution: {integrity: sha512-9h0Gvw92EvPdE8AK8AgZPbMnH5ftDyPtKm7/KUfcJVaPEPjwGDsJd1QV0H8esBDV4II41R/2lDWH1epBqIoKUw==} + engines: {node: '>=18'} + '@google-cloud/projectify@4.0.0': resolution: {integrity: sha512-MmaX6HeSvyPbWGwFq7mXdo0uQZLGBYCwziiLIGq5JVX+/bdI3SAq6bP98trV5eTWfLuvsMcIC1YJOF2vfteLFA==} engines: {node: '>=14.0.0'} @@ -973,6 +992,10 @@ packages: resolution: {integrity: sha512-Orxzlfb9c67A15cq2JQEyVc7wEsmFBmHjZWZYQMUyJ1qivXyMwdyNOs9odi79hze+2zqdTtu1E19IM/FtqZ10g==} engines: {node: '>=14'} + '@google-cloud/promisify@5.0.0': + resolution: {integrity: sha512-N8qS6dlORGHwk7WjGXKOSsLjIjNINCPicsOX6gyyLiYk7mq3MtII96NZ9N2ahwA2vnkLmZODOIH9rlNniYWvCQ==} + engines: {node: '>=18'} + '@google-cloud/storage@7.16.0': resolution: {integrity: sha512-7/5LRgykyOfQENcm6hDKP8SX/u9XxE5YOiWOkgkwcoO+cG8xT/cyOvp9wwN3IxfdYgpHs8CE7Nq2PKX2lNaEXw==} engines: {node: '>=14'} @@ -3622,6 +3645,10 @@ packages: resolution: {integrity: sha512-3duEwti880xqi4eAMN8AyR4a0ByT90zoYdLlevfrvU43vb0YZwZVfxOgxWrLXXXpyugL0hNZc9G6BiB5B3nUug==} engines: {node: '>=8'} + arrify@3.0.0: + resolution: {integrity: sha512-tLkvA81vQG/XqE2mjDkGQHoOINtMHtysSnemrmoGe6PydDPMRbVugqyk4A6V/WDWEfm3l+0d8anA9r8cv/5Jaw==} + engines: {node: '>=12'} + asap@2.0.6: resolution: {integrity: sha512-BSHWgDSAiKs50o2Re8ppvp3seVHXSRM44cdSsT9FfNEUUZLOGWVCsiWaRPWM1Znn+mqZ1OfVZ3z3DWEzSp7hRA==} @@ -4028,6 +4055,10 @@ packages: resolution: {integrity: sha512-9+vem03dMXG7gDmZ62uqmRiMRNtinIZ9ZyuF6BdxzfOD+FdN5hretzynkn0ReS2DO2GSw76RWHs0UmJPI2zUjw==} engines: {node: '>=18'} + data-uri-to-buffer@4.0.1: + resolution: {integrity: sha512-0R9ikRb668HB7QDxT1vkpuUBtqc53YyAwMwGeUFKRojY/NWKvdZ+9UYtRfGmhqNbRkTSVpMbmyhXipFFv2cb/A==} + engines: {node: '>= 12'} + data-urls@5.0.0: resolution: {integrity: sha512-ZYP5VBHshaDAiVZxjbRVcFJpc+4xGgT0bK3vzy1HLN8jTO975HEbuYzZJcHoQEY5K1a0z8YayJkyVETa08eNTg==} engines: {node: '>=18'} @@ -4451,6 +4482,10 @@ packages: fecha@4.2.3: resolution: {integrity: sha512-OP2IUU6HeYKJi3i0z4A19kHMQoLVs4Hc+DPqqxI2h/DPZHTm/vjsfC6P0b4jCMy14XizLBqvndQ+UilD7707Jw==} + fetch-blob@3.2.0: + resolution: {integrity: sha512-7yAQpD2UMJzLi1Dqv7qFYnPbaPx7ZfFK6PiIxQ4PfkGPyNyl2Ugx+a/umUonmKqjhM4DnfbMvdX6otXq83soQQ==} + engines: {node: ^12.20 || >= 14.13} + file-uri-to-path@1.0.0: resolution: {integrity: sha512-0Zt+s3L7Vf1biwWZ29aARiVYLx7iMGnEUl9x33fbB/j3jR81u/O2LbqK+Bm1CDSNDKVtJ/YjwY7TUd5SkeLQLw==} @@ -4510,6 +4545,10 @@ packages: engines: {node: '>=18.3.0'} hasBin: true + formdata-polyfill@4.0.10: + resolution: {integrity: sha512-buewHzMvYL29jdeQTVILecSaZKnt/RJWjoZCF5OW60Z67/GmSLBkOFM7qh1PI3zFNtJbaZL5eQu1vLfazOwj4g==} + engines: {node: '>=12.20.0'} + formidable@2.1.5: resolution: {integrity: sha512-Oz5Hwvwak/DCaXVVUtPn4oLMLLy1CdclLKO1LFgU7XzDpVMUU5UjlSLpGMocyQNNk8F6IJW9M/YdooSn2MRI+Q==} @@ -4542,10 +4581,18 @@ packages: resolution: {integrity: sha512-LDODD4TMYx7XXdpwxAVRAIAuB0bzv0s+ywFonY46k126qzQHT9ygyoa9tncmOiQmmDrik65UYsEkv3lbfqQ3yQ==} engines: {node: '>=14'} + gaxios@7.1.2: + resolution: {integrity: sha512-/Szrn8nr+2TsQT1Gp8iIe/BEytJmbyfrbFh419DfGQSkEgNEhbPi7JRJuughjkTzPWgU9gBQf5AVu3DbHt0OXA==} + engines: {node: '>=18'} + gcp-metadata@6.1.1: resolution: {integrity: sha512-a4tiq7E0/5fTjxPAaH4jpjkSv/uCaU2p5KC6HVGrvl0cDjA8iBZv4vv1gyzlmK0ZUKqwpOyQMKzZQe3lTit77A==} engines: {node: '>=14'} + gcp-metadata@7.0.1: + resolution: {integrity: sha512-UcO3kefx6dCcZkgcTGgVOTFb7b1LlQ02hY1omMjjrrBzkajRMCFgYOjs7J71WqnuG1k2b+9ppGL7FsOfhZMQKQ==} + engines: {node: '>=18'} + gensync@1.0.0-beta.2: resolution: {integrity: sha512-3hN7NaskYvMDLQY55gnW3NQ+mesEAepTqlg+VEbj7zzqEMBVNhzcGYYeqFo/TlYz6eQiFcp1HcsCZO+nGgS8zg==} engines: {node: '>=6.9.0'} @@ -4606,6 +4653,10 @@ packages: resolution: {integrity: sha512-WOBp/EEGUiIsJSp7wcv/y6MO+lV9UoncWqxuFfm8eBwzWNgyfBd6Gz+IeKQ9jCmyhoH99g15M3T+QaVHFjizVA==} engines: {node: '>=4'} + google-auth-library@10.3.0: + resolution: {integrity: sha512-ylSE3RlCRZfZB56PFJSfUCuiuPq83Fx8hqu1KPWGK8FVdSaxlp/qkeMMX/DT/18xkwXIHvXEXkZsljRwfrdEfQ==} + engines: {node: '>=18'} + google-auth-library@9.15.1: resolution: {integrity: sha512-Jb6Z0+nvECVz+2lzSMt9u98UsoakXxA2HGHMCxh+so3n90XgYWkq5dur19JAJV7ONiJY22yBTyJB1TSkvPq9Ng==} engines: {node: '>=14'} @@ -4614,6 +4665,10 @@ packages: resolution: {integrity: sha512-NEgUnEcBiP5HrPzufUkBzJOD/Sxsco3rLNo1F1TNf7ieU8ryUzBhqba8r756CjLX7rn3fHl6iLEwPYuqpoKgQQ==} engines: {node: '>=14'} + google-logging-utils@1.1.1: + resolution: {integrity: sha512-rcX58I7nqpu4mbKztFeOAObbomBbHU2oIb/d3tJfF3dizGSApqtSwYJigGCooHdnMyQBIw8BrWyK96w3YXgr6A==} + engines: {node: '>=14'} + gopd@1.0.1: resolution: {integrity: sha512-d65bNlIadxvpb/A2abVdlqKqV563juRnZ1Wtk6s1sIR8uNsXR70xqIzVqxVf1eTqDunwT2MkczEeaezCKTZhwA==} @@ -4628,6 +4683,10 @@ packages: resolution: {integrity: sha512-pCcEwRi+TKpMlxAQObHDQ56KawURgyAf6jtIY046fJ5tIv3zDe/LEIubckAO8fj6JnAxLdmWkUfNyulQ2iKdEw==} engines: {node: '>=14.0.0'} + gtoken@8.0.0: + resolution: {integrity: sha512-+CqsMbHPiSTdtSO14O51eMNlrp9N79gmeqmXeouJOhfucAedHw9noVe/n5uJk3tbKE6a+6ZCQg3RPhVhHByAIw==} + engines: {node: '>=18'} + h3@1.15.4: resolution: {integrity: sha512-z5cFQWDffyOe4vQ9xIqNfCZdV4p//vy6fBnr8Q1AWnVZ0teurKMG66rLj++TKwKPUP3u7iMUvrvKaEUiQw2QWQ==} @@ -5469,6 +5528,11 @@ packages: node-cleanup@2.1.2: resolution: {integrity: sha512-qN8v/s2PAJwGUtr1/hYTpNKlD6Y9rc4p8KSmJXyGdYGZsDGKXrGThikLFP9OCHFeLeEpQzPwiAtdIvBLqm//Hw==} + node-domexception@1.0.0: + resolution: {integrity: sha512-/jKZoMpw0F8GRwl4/eLROPA3cfcXtLApP0QzLmUT/HuPCZWyB7IY9ZrMeKw2O/nFIqPQB3PVM9aYm0F312AXDQ==} + engines: {node: '>=10.5.0'} + deprecated: Use your platform's native DOMException instead + node-ensure@0.0.0: resolution: {integrity: sha512-DRI60hzo2oKN1ma0ckc6nQWlHU69RH6xN0sjQTjMpChPfTYvKZdcQFfdYK2RWbJcKyUizSIy/l8OTGxMAM1QDw==} @@ -5484,6 +5548,10 @@ packages: encoding: optional: true + node-fetch@3.3.2: + resolution: {integrity: sha512-dRB78srN/l6gqWulah9SrxeYnxeddIG30+GOqK/9OlLVyLg3HPnr6SqOWTWOXKRwC2eGYCkZ59NNuSgvSrpgOA==} + engines: {node: ^12.20.0 || ^14.13.1 || >=16.0.0} + node-gyp-build-optional-packages@5.2.2: resolution: {integrity: sha512-s+w+rBWnpTMwSFbaE0UXsRlg7hU4FjekKU4eyAih5T8nJuNZT1nNsskXpxmeqSK9UzkBl6UgRlnKc8hz8IEqOw==} hasBin: true @@ -6057,6 +6125,10 @@ packages: resolution: {integrity: sha512-dUOvLMJ0/JJYEn8NrpOaGNE7X3vpI5XlZS/u0ANjqtcZVKnIxP7IgCFwrKTxENw29emmwug53awKtaMm4i9g5w==} engines: {node: '>=14'} + retry-request@8.0.2: + resolution: {integrity: sha512-JzFPAfklk1kjR1w76f0QOIhoDkNkSqW8wYKT08n9yysTmZfB+RQ2QoXoTAeOi1HD9ZipTyTAZg3c4pM/jeqgSw==} + engines: {node: '>=18'} + retry@0.13.1: resolution: {integrity: sha512-XQBQ3I8W1Cge0Seh+6gjj03LbmRFWuoszgK9ooCpwYIrhhoO80pfq4cUkU5DkknwfOfFteRwlZ56PYOGYyFWdg==} engines: {node: '>= 4'} @@ -6403,6 +6475,10 @@ packages: os: [darwin, linux, win32, freebsd, openbsd, netbsd, sunos, android] hasBin: true + teeny-request@10.1.0: + resolution: {integrity: sha512-3ZnLvgWF29jikg1sAQ1g0o+lr5JX6sVgYvfUJazn7ZjJroDBUTWp44/+cFVX0bULjv4vci+rBD+oGVAkWqhUbw==} + engines: {node: '>=18'} + teeny-request@9.0.0: resolution: {integrity: sha512-resvxdc6Mgb7YEThw6G6bExlXKkv6+YbuzGg9xuXxSgxJF7Ozs+o8Y9+2R3sArdWdW8nOokoQb1yrpFB0pQK2g==} engines: {node: '>=14'} @@ -6770,6 +6846,10 @@ packages: walker@1.0.8: resolution: {integrity: sha512-ts/8E8l5b7kY0vlWLewOkDXMmPdLcVV4GmOQLyxuSswIJsweeFZtAsMF7k1Nszz+TYBQrlYRmzOnr398y1JemQ==} + web-streams-polyfill@3.3.3: + resolution: {integrity: sha512-d2JWLCivmZYTSIoge9MsgFCZrt571BikcWGYkjC1khllbTeDlGqZ2D8vD8E/lJa8WGWbb7Plm8/XJYV7IJHZZw==} + engines: {node: '>= 8'} + webextension-polyfill@0.10.0: resolution: {integrity: sha512-c5s35LgVa5tFaHhrZDnr3FpQpjj1BB+RXhLTYUxGqBVN460HkbM8TBtEqdXWbpTKfzwCcjAZVF7zXCYSKtcp9g==} @@ -7733,15 +7813,52 @@ snapshots: ethereum-cryptography: 2.2.1 micro-ftch: 0.3.1 + '@google-cloud/bigquery@8.1.1': + dependencies: + '@google-cloud/common': 6.0.0 + '@google-cloud/paginator': 6.0.0 + '@google-cloud/precise-date': 5.0.0 + '@google-cloud/promisify': 5.0.0 + arrify: 3.0.0 + big.js: 6.2.2 + duplexify: 4.1.3 + extend: 3.0.2 + stream-events: 1.0.5 + teeny-request: 10.1.0 + transitivePeerDependencies: + - supports-color + + '@google-cloud/common@6.0.0': + dependencies: + '@google-cloud/projectify': 4.0.0 + '@google-cloud/promisify': 4.0.0 + arrify: 2.0.1 + duplexify: 4.1.3 + extend: 3.0.2 + google-auth-library: 10.3.0 + html-entities: 2.6.0 + retry-request: 8.0.2 + teeny-request: 10.1.0 + transitivePeerDependencies: + - supports-color + '@google-cloud/paginator@5.0.2': dependencies: arrify: 2.0.1 extend: 3.0.2 + '@google-cloud/paginator@6.0.0': + dependencies: + extend: 3.0.2 + + '@google-cloud/precise-date@5.0.0': {} + '@google-cloud/projectify@4.0.0': {} '@google-cloud/promisify@4.0.0': {} + '@google-cloud/promisify@5.0.0': {} + '@google-cloud/storage@7.16.0(encoding@0.1.13)': dependencies: '@google-cloud/paginator': 5.0.2 @@ -11672,6 +11789,8 @@ snapshots: arrify@2.0.1: {} + arrify@3.0.0: {} + asap@2.0.6: {} async-mutex@0.2.6: @@ -12143,6 +12262,8 @@ snapshots: '@asamuzakjp/css-color': 2.8.3 rrweb-cssom: 0.8.0 + data-uri-to-buffer@4.0.1: {} + data-urls@5.0.0: dependencies: whatwg-mimetype: 4.0.0 @@ -12582,6 +12703,11 @@ snapshots: fecha@4.2.3: {} + fetch-blob@3.2.0: + dependencies: + node-domexception: 1.0.0 + web-streams-polyfill: 3.3.3 + file-uri-to-path@1.0.0: {} filelist@1.0.4: @@ -12651,6 +12777,10 @@ snapshots: dependencies: fd-package-json: 2.0.0 + formdata-polyfill@4.0.10: + dependencies: + fetch-blob: 3.2.0 + formidable@2.1.5: dependencies: '@paralleldrive/cuid2': 2.2.2 @@ -12684,6 +12814,14 @@ snapshots: - encoding - supports-color + gaxios@7.1.2: + dependencies: + extend: 3.0.2 + https-proxy-agent: 7.0.6 + node-fetch: 3.3.2 + transitivePeerDependencies: + - supports-color + gcp-metadata@6.1.1(encoding@0.1.13): dependencies: gaxios: 6.7.1(encoding@0.1.13) @@ -12693,6 +12831,14 @@ snapshots: - encoding - supports-color + gcp-metadata@7.0.1: + dependencies: + gaxios: 7.1.2 + google-logging-utils: 1.1.1 + json-bigint: 1.0.0 + transitivePeerDependencies: + - supports-color + gensync@1.0.0-beta.2: {} geoip-country@5.0.202508192342: @@ -12772,6 +12918,18 @@ snapshots: globals@11.12.0: {} + google-auth-library@10.3.0: + dependencies: + base64-js: 1.5.1 + ecdsa-sig-formatter: 1.0.11 + gaxios: 7.1.2 + gcp-metadata: 7.0.1 + google-logging-utils: 1.1.1 + gtoken: 8.0.0 + jws: 4.0.0 + transitivePeerDependencies: + - supports-color + google-auth-library@9.15.1(encoding@0.1.13): dependencies: base64-js: 1.5.1 @@ -12786,6 +12944,8 @@ snapshots: google-logging-utils@0.0.2: {} + google-logging-utils@1.1.1: {} + gopd@1.0.1: dependencies: get-intrinsic: 1.3.0 @@ -12802,6 +12962,13 @@ snapshots: - encoding - supports-color + gtoken@8.0.0: + dependencies: + gaxios: 7.1.2 + jws: 4.0.0 + transitivePeerDependencies: + - supports-color + h3@1.15.4: dependencies: cookie-es: 1.2.2 @@ -13849,6 +14016,8 @@ snapshots: node-cleanup@2.1.2: {} + node-domexception@1.0.0: {} + node-ensure@0.0.0: {} node-fetch-native@1.6.7: {} @@ -13859,6 +14028,12 @@ snapshots: optionalDependencies: encoding: 0.1.13 + node-fetch@3.3.2: + dependencies: + data-uri-to-buffer: 4.0.1 + fetch-blob: 3.2.0 + formdata-polyfill: 4.0.10 + node-gyp-build-optional-packages@5.2.2: dependencies: detect-libc: 2.0.3 @@ -14492,6 +14667,13 @@ snapshots: - encoding - supports-color + retry-request@8.0.2: + dependencies: + extend: 3.0.2 + teeny-request: 10.1.0 + transitivePeerDependencies: + - supports-color + retry@0.13.1: {} reusify@1.1.0: {} @@ -14857,6 +15039,15 @@ snapshots: systeminformation@5.27.8: {} + teeny-request@10.1.0: + dependencies: + http-proxy-agent: 5.0.0 + https-proxy-agent: 5.0.1 + node-fetch: 3.3.2 + stream-events: 1.0.5 + transitivePeerDependencies: + - supports-color + teeny-request@9.0.0(encoding@0.1.13): dependencies: http-proxy-agent: 5.0.0 @@ -15228,6 +15419,8 @@ snapshots: dependencies: makeerror: 1.0.12 + web-streams-polyfill@3.3.3: {} + webextension-polyfill@0.10.0: {} webidl-conversions@3.0.1: {} diff --git a/apps/api/src/lib/bigquery-jobs.ts b/apps/api/src/lib/bigquery-jobs.ts new file mode 100644 index 000000000..f117110dd --- /dev/null +++ b/apps/api/src/lib/bigquery-jobs.ts @@ -0,0 +1,370 @@ +import { FirecrawlJob } from "../types"; +import { logger } from "./logger"; +import { transformJobForLogging, createJobLoggerContext } from "./job-transform"; + +// BigQuery client will be imported conditionally to avoid errors if not installed +let BigQuery: any = null; +let bigquery: any = null; + +try { + // Dynamically import BigQuery to handle cases where it's not installed + BigQuery = require("@google-cloud/bigquery").BigQuery; + + const credentials = process.env.GCS_CREDENTIALS + ? JSON.parse(atob(process.env.GCS_CREDENTIALS)) + : undefined; + + bigquery = new BigQuery({ + projectId: process.env.GOOGLE_CLOUD_PROJECT_ID || "firecrawl", + credentials, + }); +} catch (error) { + logger.warn("BigQuery client not available. BigQuery logging will be disabled.", { + error: error.message, + }); +} + +/** + * Transforms a FirecrawlJob into a BigQuery-compatible row + */ +function transformJobForBigQuery(job: FirecrawlJob) { + const transformed = transformJobForLogging(job, { + includeTimestamp: true, + serializeObjects: true, + cleanNullValues: false, // BigQuery handles null values fine + }); + + // Remove docs field from BigQuery to reduce storage costs and avoid size issues + const { docs, ...bigQueryRow } = transformed; + + return bigQueryRow; +} + +/** + * Ensures the BigQuery dataset exists + */ +async function ensureBigQueryDataset(): Promise { + if (!bigquery || !process.env.BIGQUERY_DATASET_ID) { + return; + } + + try { + const datasetId = process.env.BIGQUERY_DATASET_ID; + const dataset = bigquery.dataset(datasetId); + + const [exists] = await dataset.exists(); + if (!exists) { + logger.info("Creating BigQuery dataset", { dataset: datasetId }); + + await dataset.create({ + location: process.env.BIGQUERY_LOCATION || "US", + }); + + logger.info("BigQuery dataset created successfully", { dataset: datasetId }); + } + } catch (error) { + logger.error("Error ensuring BigQuery dataset exists", { + error, + dataset: process.env.BIGQUERY_DATASET_ID, + }); + throw error; + } +} + +/** + * Creates the BigQuery table if it doesn't exist + */ +async function ensureBigQueryTable(): Promise { + if (!bigquery || !process.env.BIGQUERY_DATASET_ID) { + return; + } + + try { + // Ensure dataset exists first + // No need once it is initialized + // await ensureBigQueryDataset(); + + const datasetId = process.env.BIGQUERY_DATASET_ID; + const tableId = process.env.BIGQUERY_TABLE_ID || "firecrawl_jobs"; + + const dataset = bigquery.dataset(datasetId); + const table = dataset.table(tableId); + + const [exists] = await table.exists(); + + if (exists) { + // Table exists, but let's check if it has the correct schema (without docs field) + try { + const [metadata] = await table.getMetadata(); + const currentSchema = metadata.schema; + + if (!currentSchema || !currentSchema.fields || currentSchema.fields.length === 0) { + logger.warn("Table exists but has no schema, recreating table", { + dataset: datasetId, + table: tableId, + }); + + // Delete and recreate the table + await table.delete(); + logger.info("Deleted table with no schema", { dataset: datasetId, table: tableId }); + } else { + // Check if the table has the old schema with docs field or wrong data types + const hasDocsField = currentSchema.fields.some((field: any) => field.name === 'docs'); + const timeTakenField = currentSchema.fields.find((field: any) => field.name === 'time_taken'); + const numTokensField = currentSchema.fields.find((field: any) => field.name === 'num_tokens'); + const hasWrongTimeTakenType = timeTakenField && timeTakenField.type !== 'FLOAT'; + const hasWrongNumTokensType = numTokensField && numTokensField.type !== 'FLOAT'; + + if (hasDocsField) { + logger.warn("Table has old schema with docs field, recreating table", { + dataset: datasetId, + table: tableId, + }); + + // Delete and recreate the table without docs field + await table.delete(); + logger.info("Deleted table with old schema (docs field)", { dataset: datasetId, table: tableId }); + } else if (hasWrongTimeTakenType || hasWrongNumTokensType) { + logger.warn("Table has wrong data types for numeric fields, recreating table", { + dataset: datasetId, + table: tableId, + timeTakenType: timeTakenField?.type, + numTokensType: numTokensField?.type, + expectedType: 'FLOAT', + }); + + // Delete and recreate the table with correct data types + await table.delete(); + logger.info("Deleted table with wrong schema (numeric field types)", { dataset: datasetId, table: tableId }); + } else { + logger.debug("BigQuery table exists and has correct schema", { + dataset: datasetId, + table: tableId, + }); + return; + } + } + } catch (schemaError) { + logger.error("Error checking table schema, recreating table", { + error: schemaError, + dataset: datasetId, + table: tableId, + }); + + // Try to delete and recreate + try { + await table.delete(); + logger.info("Deleted problematic table", { dataset: datasetId, table: tableId }); + } catch (deleteError) { + logger.error("Error deleting problematic table", { error: deleteError }); + } + } + } + + // Create the table (either it doesn't exist or we deleted it) + logger.info("Creating BigQuery table", { + dataset: datasetId, + table: tableId, + }); + + const schema = [ + { name: "job_id", type: "STRING", mode: "NULLABLE" }, + { name: "success", type: "BOOLEAN", mode: "NULLABLE" }, + { name: "message", type: "STRING", mode: "NULLABLE" }, + { name: "num_docs", type: "INTEGER", mode: "NULLABLE" }, + // Note: docs field excluded from BigQuery to reduce storage costs and avoid size issues + { name: "time_taken", type: "FLOAT", mode: "NULLABLE" }, + { name: "team_id", type: "STRING", mode: "NULLABLE" }, + { name: "mode", type: "STRING", mode: "NULLABLE" }, + { name: "url", type: "STRING", mode: "NULLABLE" }, + { name: "crawler_options", type: "STRING", mode: "NULLABLE" }, + { name: "page_options", type: "STRING", mode: "NULLABLE" }, + { name: "origin", type: "STRING", mode: "NULLABLE" }, + { name: "integration", type: "STRING", mode: "NULLABLE" }, + { name: "num_tokens", type: "FLOAT", mode: "NULLABLE" }, + { name: "retry", type: "BOOLEAN", mode: "NULLABLE" }, + { name: "crawl_id", type: "STRING", mode: "NULLABLE" }, + { name: "tokens_billed", type: "INTEGER", mode: "NULLABLE" }, + { name: "is_migrated", type: "BOOLEAN", mode: "NULLABLE" }, + { name: "cost_tracking", type: "STRING", mode: "NULLABLE" }, + { name: "pdf_num_pages", type: "INTEGER", mode: "NULLABLE" }, + { name: "credits_billed", type: "INTEGER", mode: "NULLABLE" }, + { name: "change_tracking_tag", type: "STRING", mode: "NULLABLE" }, + { name: "dr_clean_by", type: "TIMESTAMP", mode: "NULLABLE" }, + { name: "timestamp", type: "TIMESTAMP", mode: "NULLABLE" }, + ]; + + const options = { + schema: { fields: schema }, + timePartitioning: { + type: "HOUR", + field: "timestamp", + }, + location: process.env.BIGQUERY_LOCATION || "US", + }; + + await table.create(options); + + logger.info("BigQuery table created successfully", { + dataset: datasetId, + table: tableId, + schema: schema.length + " fields", + }); + } catch (error) { + logger.error("Error ensuring BigQuery table exists", { + error, + dataset: process.env.BIGQUERY_DATASET_ID, + table: process.env.BIGQUERY_TABLE_ID || "firecrawl_jobs", + }); + throw error; + } +} + +/** + * Validates BigQuery configuration + */ +function validateBigQueryConfig(): { valid: boolean; error?: string } { + if (!bigquery) { + return { valid: false, error: "BigQuery client not initialized" }; + } + + if (!process.env.BIGQUERY_DATASET_ID) { + return { valid: false, error: "BIGQUERY_DATASET_ID not configured" }; + } + + if (!process.env.GOOGLE_CLOUD_PROJECT_ID && !process.env.GCS_CREDENTIALS) { + return { valid: false, error: "Google Cloud credentials not configured" }; + } + + return { valid: true }; +} + +/** + * Saves a job to BigQuery + */ +export async function saveJobToBigQuery( + job: FirecrawlJob, + force: boolean = false +): Promise { + const configValidation = validateBigQueryConfig(); + if (!configValidation.valid) { + logger.debug("BigQuery not configured, skipping BigQuery logging", { + reason: configValidation.error, + }); + return; + } + + const jobLogger = logger.child({ + module: "bigquery_jobs", + method: "saveJobToBigQuery", + ...createJobLoggerContext(job), + }); + + try { + // Ensure table exists before attempting to insert + await ensureBigQueryTable(); + + const datasetId = process.env.BIGQUERY_DATASET_ID; + const tableId = process.env.BIGQUERY_TABLE_ID || "firecrawl_jobs"; + const table = bigquery.dataset(datasetId).table(tableId); + + const row = transformJobForBigQuery(job) as any; + + // Validate the row has required fields + if (!row.timestamp) { + row.timestamp = new Date().toISOString(); + } + + jobLogger.debug("Attempting to insert job to BigQuery", { + dataset: datasetId, + table: tableId, + jobId: job.job_id, + }); + + if (force) { + let i = 0; + let done = false; + while (i++ <= 10 && !done) { + try { + await table.insert([row], { + ignoreUnknownValues: false, + skipInvalidRows: false, + }); + done = true; + jobLogger.debug("Job logged to BigQuery successfully!"); + } catch (error) { + jobLogger.error( + "Failed to log job to BigQuery due to error -- trying again", + { + error: error.message || error, + attempt: i, + jobId: job.job_id, + dataset: datasetId, + table: tableId, + } + ); + + // If it's a schema error, don't retry + if (error.message && error.message.includes("schema")) { + jobLogger.error("Schema error detected, stopping retries", { error: error.message }); + break; + } + + await new Promise((resolve) => setTimeout(() => resolve(), 100 * i)); + } + } + if (!done) { + jobLogger.error("Failed to log job to BigQuery after all retries!"); + } + } else { + try { + await table.insert([row], { + ignoreUnknownValues: false, + skipInvalidRows: false, + }); + jobLogger.debug("Job logged to BigQuery successfully!"); + } catch (error) { + jobLogger.error("Error logging job to BigQuery", { + error: error.message || error, + jobId: job.job_id, + dataset: datasetId, + table: tableId, + }); + } + } + } catch (error) { + jobLogger.error("Error saving job to BigQuery", { + error: error.message || error, + jobId: job.job_id, + dataset: process.env.BIGQUERY_DATASET_ID, + table: process.env.BIGQUERY_TABLE_ID || "firecrawl_jobs", + }); + } +} + +/** + * Queries jobs from BigQuery + */ +export async function queryJobsFromBigQuery( + query: string, + params?: any[] +): Promise { + if (!bigquery || !process.env.BIGQUERY_DATASET_ID) { + logger.warn("BigQuery not configured"); + return []; + } + + try { + const options = { + query, + params, + location: process.env.BIGQUERY_LOCATION || "US", + }; + + const [rows] = await bigquery.query(options); + return rows; + } catch (error) { + logger.error("Error querying jobs from BigQuery", { error, query }); + throw error; + } +} diff --git a/apps/api/src/lib/job-transform.ts b/apps/api/src/lib/job-transform.ts new file mode 100644 index 000000000..38746ef32 --- /dev/null +++ b/apps/api/src/lib/job-transform.ts @@ -0,0 +1,131 @@ +import { FirecrawlJob } from "../types"; + +function cleanOfNull(x: T): T { + if (Array.isArray(x)) { + return x.map(x => cleanOfNull(x)) as T; + } else if (typeof x === "object" && x !== null) { + return Object.fromEntries( + Object.entries(x).map(([k, v]) => [k, cleanOfNull(v)]), + ) as T; + } else if (typeof x === "string") { + return x.replaceAll("\u0000", "") as T; + } else { + return x; + } +} + +export interface TransformOptions { + /** Whether to include timestamp field (for BigQuery) */ + includeTimestamp?: boolean; + /** Whether to serialize objects to JSON strings (for BigQuery) */ + serializeObjects?: boolean; + /** Whether to clean null values from docs */ + cleanNullValues?: boolean; +} + +/** + * Transforms a FirecrawlJob into a standardized format for logging/storage + */ +export function transformJobForLogging( + job: FirecrawlJob, + options: TransformOptions = {} +) { + const { + includeTimestamp = false, + serializeObjects = false, + cleanNullValues = true, + } = options; + + const zeroDataRetention = job.zeroDataRetention ?? false; + + // Determine if docs should be included based on zero data retention and GCS usage + const shouldIncludeDocs = !zeroDataRetention && + !((job.mode === "single_urls" || job.mode === "scrape") && process.env.GCS_BUCKET_NAME); + + const baseTransform = { + job_id: job.job_id ? job.job_id : null, + success: job.success, + message: zeroDataRetention ? null : job.message, + num_docs: job.num_docs, + docs: shouldIncludeDocs + ? (cleanNullValues ? cleanOfNull(job.docs) : job.docs) + : null, + time_taken: job.time_taken, + team_id: + job.team_id === "preview" || job.team_id?.startsWith("preview_") + ? null + : job.team_id, + mode: job.mode, + url: zeroDataRetention + ? "" + : job.url, + crawler_options: zeroDataRetention ? null : job.crawlerOptions, + page_options: zeroDataRetention ? null : job.scrapeOptions, + origin: zeroDataRetention ? null : job.origin, + integration: zeroDataRetention ? null : (job.integration ?? null), + num_tokens: job.num_tokens, + retry: !!job.retry, + crawl_id: job.crawl_id, + tokens_billed: job.tokens_billed, + is_migrated: true, + cost_tracking: zeroDataRetention ? null : job.cost_tracking, + pdf_num_pages: zeroDataRetention ? null : (job.pdf_num_pages ?? null), + credits_billed: job.credits_billed ?? null, + change_tracking_tag: zeroDataRetention + ? null + : (job.change_tracking_tag ?? null), + dr_clean_by: + zeroDataRetention && job.crawl_id + ? new Date(Date.now() + 1000 * 60 * 60 * 24).toISOString() + : null, + }; + + // Add timestamp if requested (for BigQuery) + const withTimestamp = includeTimestamp + ? { ...baseTransform, timestamp: new Date().toISOString() } + : baseTransform; + + // Serialize objects to JSON strings if requested (for BigQuery) + if (serializeObjects) { + return { + ...withTimestamp, + docs: withTimestamp.docs ? JSON.stringify(withTimestamp.docs) : null, + crawler_options: withTimestamp.crawler_options + ? JSON.stringify(withTimestamp.crawler_options) + : null, + page_options: withTimestamp.page_options + ? JSON.stringify(withTimestamp.page_options) + : null, + cost_tracking: withTimestamp.cost_tracking + ? JSON.stringify(withTimestamp.cost_tracking) + : null, + }; + } + + return withTimestamp; +} + +/** + * Creates logger context based on job mode and ID + */ +export function createJobLoggerContext(job: FirecrawlJob) { + return { + ...(job.mode === "scrape" || + job.mode === "single_urls" || + job.mode === "single_url" + ? { + scrapeId: job.job_id, + } + : {}), + ...(job.mode === "crawl" || job.mode === "batch_scrape" + ? { + crawlId: job.job_id, + } + : {}), + ...(job.mode === "extract" + ? { + extractId: job.job_id, + } + : {}), + }; +} diff --git a/apps/api/src/services/logging/log_job.ts b/apps/api/src/services/logging/log_job.ts index cd0907205..9d9e23faf 100644 --- a/apps/api/src/services/logging/log_job.ts +++ b/apps/api/src/services/logging/log_job.ts @@ -5,21 +5,10 @@ import "dotenv/config"; import { logger as _logger } from "../../lib/logger"; import { configDotenv } from "dotenv"; import { saveJobToGCS } from "../../lib/gcs-jobs"; +import { saveJobToBigQuery } from "../../lib/bigquery-jobs"; +import { transformJobForLogging, createJobLoggerContext } from "../../lib/job-transform"; configDotenv(); -function cleanOfNull(x: T): T { - if (Array.isArray(x)) { - return x.map(x => cleanOfNull(x)) as T; - } else if (typeof x === "object" && x !== null) { - return Object.fromEntries( - Object.entries(x).map(([k, v]) => [k, cleanOfNull(v)]), - ) as T; - } else if (typeof x === "string") { - return x.replaceAll("\u0000", "") as T; - } else { - return x; - } -} export async function logJob( job: FirecrawlJob, @@ -29,23 +18,7 @@ export async function logJob( let logger = _logger.child({ module: "log_job", method: "logJob", - ...(job.mode === "scrape" || - job.mode === "single_urls" || - job.mode === "single_url" - ? { - scrapeId: job.job_id, - } - : {}), - ...(job.mode === "crawl" || job.mode === "batch_scrape" - ? { - crawlId: job.job_id, - } - : {}), - ...(job.mode === "extract" - ? { - extractId: job.job_id, - } - : {}), + ...createJobLoggerContext(job), }); const zeroDataRetention = job.zeroDataRetention ?? false; @@ -55,10 +28,18 @@ export async function logJob( }); try { + // Save to GCS if configured if (process.env.GCS_BUCKET_NAME) { await saveJobToGCS(job); } + // Save to BigQuery if configured + if (process.env.BIGQUERY_DATASET_ID) { + saveJobToBigQuery(job, force).catch(error => { + logger.error("Error saving job to BigQuery", { error }); + }); + } + const useDbAuthentication = process.env.USE_DB_AUTHENTICATION === "true"; if (!useDbAuthentication) { return; @@ -79,46 +60,11 @@ export async function logJob( // }, // ]; // } - const jobColumn = { - job_id: job.job_id ? job.job_id : null, - success: job.success, - message: zeroDataRetention ? null : job.message, - num_docs: job.num_docs, - docs: zeroDataRetention - ? null - : (job.mode === "single_urls" || job.mode === "scrape") && - process.env.GCS_BUCKET_NAME - ? null - : cleanOfNull(job.docs), - time_taken: job.time_taken, - team_id: - job.team_id === "preview" || job.team_id?.startsWith("preview_") - ? null - : job.team_id, - mode: job.mode, - url: zeroDataRetention - ? "" - : job.url, - crawler_options: zeroDataRetention ? null : job.crawlerOptions, - page_options: zeroDataRetention ? null : job.scrapeOptions, - origin: zeroDataRetention ? null : job.origin, - integration: zeroDataRetention ? null : (job.integration ?? null), - num_tokens: job.num_tokens, - retry: !!job.retry, - crawl_id: job.crawl_id, - tokens_billed: job.tokens_billed, - is_migrated: true, - cost_tracking: zeroDataRetention ? null : job.cost_tracking, - pdf_num_pages: zeroDataRetention ? null : (job.pdf_num_pages ?? null), - credits_billed: job.credits_billed ?? null, - change_tracking_tag: zeroDataRetention - ? null - : (job.change_tracking_tag ?? null), - dr_clean_by: - zeroDataRetention && job.crawl_id - ? new Date(Date.now() + 1000 * 60 * 60 * 24).toISOString() - : null, - }; + const jobColumn = transformJobForLogging(job, { + includeTimestamp: false, + serializeObjects: false, + cleanNullValues: true, + }); if (bypassLogging) { return; @@ -169,6 +115,12 @@ export async function logJob( } if (process.env.POSTHOG_API_KEY && !job.crawl_id) { + const jobProperties = transformJobForLogging(job, { + includeTimestamp: false, + serializeObjects: false, + cleanNullValues: false, + }); + let phLog = { distinctId: "from-api", //* To identify this on the group level, setting distinctid to a static string per posthog docs: https://posthog.com/docs/product-analytics/group-analytics#advanced-server-side-only-capturing-group-events-without-a-user ...(job.team_id !== "preview" && @@ -177,29 +129,9 @@ export async function logJob( }), //* Identifying event on this team event: "job-logged", properties: { - success: job.success, - message: zeroDataRetention ? null : job.message, - num_docs: job.num_docs, - time_taken: job.time_taken, - team_id: - job.team_id === "preview" || job.team_id?.startsWith("preview_") - ? null - : job.team_id, - mode: job.mode, - url: zeroDataRetention - ? "" - : job.url, - crawler_options: zeroDataRetention ? null : job.crawlerOptions, - page_options: zeroDataRetention ? null : job.scrapeOptions, - origin: zeroDataRetention ? null : job.origin, - num_tokens: job.num_tokens, - retry: job.retry, - tokens_billed: job.tokens_billed, - cost_tracking: zeroDataRetention ? null : job.cost_tracking, - pdf_num_pages: zeroDataRetention ? null : job.pdf_num_pages, - change_tracking_tag: zeroDataRetention - ? null - : (job.change_tracking_tag ?? null), + ...jobProperties, + // Remove docs from PostHog as it's not needed for analytics + docs: undefined, }, }; if (job.mode !== "single_urls") {