From 646efbe3fc66272ee6f39d51eb0f470bdba1ea10 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Gerg=C5=91=20M=C3=B3ricz?= Date: Wed, 22 Oct 2025 16:50:04 +0200 Subject: [PATCH] feat(nuq): group_id, job backlogs, and group add operations (#2309) * wip * feat(nuq): bulk add operations * various fixes * various fixes --- apps/api/package.json | 4 +- apps/api/pnpm-lock.yaml | 104 ++----- apps/api/src/lib/concurrency-limit.ts | 11 +- apps/api/src/services/queue-jobs.ts | 154 ++++++++-- apps/api/src/services/worker/nuq.ts | 271 ++++++++++++++++-- apps/api/src/services/worker/scrape-worker.ts | 16 +- apps/nuq-postgres/nuq.sql | 13 + 7 files changed, 444 insertions(+), 129 deletions(-) diff --git a/apps/api/package.json b/apps/api/package.json index 173187409..0fd875f51 100644 --- a/apps/api/package.json +++ b/apps/api/package.json @@ -71,7 +71,6 @@ "@ai-sdk/fireworks": "^0.2.4", "@ai-sdk/google": "^1.2.3", "@ai-sdk/google-vertex": "^2.2.15", - "@pinecone-database/pinecone": "^6.1.2", "@ai-sdk/groq": "^1.2.1", "@ai-sdk/openai": "^1.3.12", "@apidevtools/json-schema-ref-parser": "^11.7.3", @@ -88,6 +87,7 @@ "@opentelemetry/core": "^2.1.0", "@opentelemetry/exporter-trace-otlp-http": "^0.205.0", "@opentelemetry/sdk-node": "^0.205.0", + "@pinecone-database/pinecone": "^6.1.2", "@sentry/cli": "^2.33.1", "@sentry/node": "^9.40.0", "@sentry/profiling-node": "^9.40.0", @@ -137,7 +137,7 @@ "tough-cookie": "^4.1.4", "turndown": "^7.1.3", "undici": "^7.10.0", - "uuid": "^10.0.0", + "uuid": "^13.0.0", "winston": "^3.14.2", "ws": "^8.18.0", "x402-express": "^0.6.5", diff --git a/apps/api/pnpm-lock.yaml b/apps/api/pnpm-lock.yaml index 43bf374bf..0bbbc7abb 100644 --- a/apps/api/pnpm-lock.yaml +++ b/apps/api/pnpm-lock.yaml @@ -227,8 +227,8 @@ importers: specifier: ^7.10.0 version: 7.10.0 uuid: - specifier: ^10.0.0 - version: 10.0.0 + specifier: ^13.0.0 + version: 13.0.0 winston: specifier: ^3.14.2 version: 3.14.2 @@ -6751,8 +6751,8 @@ packages: resolution: {integrity: sha512-pMZTvIkT1d+TFGvDOqodOclx0QWkkgi6Tdoa8gC8ffGAAqz9pzPTZWAybbsHHoED/ztMtkv/VoYTYyShUn81hA==} engines: {node: '>= 0.4.0'} - uuid@10.0.0: - resolution: {integrity: sha512-8XkAphELsDnEGrDxUOHB3RGvXz6TeuYSGEZBOjtTtPm2lwhGBjLgOzLHB63IUWfBpNucQjND6d3AOudO+H3RWQ==} + uuid@13.0.0: + resolution: {integrity: sha512-XQegIaBTVUjSHliKqcnFqYypAd4S+WCYt5NIeRs6w/UAry7z8Y9j5ZwRRL4kzq9U3sD6v+85er9FvkEaBpji2w==} hasBin: true uuid@8.3.2: @@ -7537,7 +7537,7 @@ snapshots: idb-keyval: 6.2.1 ox: 0.6.9(typescript@5.8.3)(zod@3.25.76) preact: 10.24.2 - viem: 2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.24.2) + viem: 2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76) zustand: 5.0.3(react@18.3.1)(use-sync-external-store@1.4.0(react@18.3.1)) transitivePeerDependencies: - '@types/react' @@ -7611,7 +7611,7 @@ snapshots: idb-keyval: 6.2.1 ox: 0.6.9(typescript@5.8.3)(zod@3.25.76) preact: 10.24.2 - viem: 2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.24.2) + viem: 2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76) zustand: 5.0.3(react@18.3.1)(use-sync-external-store@1.4.0(react@18.3.1)) transitivePeerDependencies: - '@types/react' @@ -7814,11 +7814,11 @@ snapshots: ethereum-cryptography: 2.2.1 micro-ftch: 0.3.1 - '@gemini-wallet/core@0.2.0(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76))': + '@gemini-wallet/core@0.2.0(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.24.2))': dependencies: '@metamask/rpc-errors': 7.0.2 eventemitter3: 5.0.1 - viem: 2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.24.2) + viem: 2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76) transitivePeerDependencies: - supports-color @@ -9860,7 +9860,7 @@ snapshots: dependencies: big.js: 6.2.2 dayjs: 1.11.13 - viem: 2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.24.2) + viem: 2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76) transitivePeerDependencies: - bufferutil - typescript @@ -9873,7 +9873,7 @@ snapshots: '@reown/appkit-wallet': 1.7.8(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10) '@walletconnect/universal-provider': 2.21.0(bufferutil@4.0.9)(encoding@0.1.13)(ioredis@5.6.1)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76) valtio: 1.13.2(react@18.3.1) - viem: 2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.24.2) + viem: 2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76) transitivePeerDependencies: - '@azure/app-configuration' - '@azure/cosmos' @@ -10019,7 +10019,7 @@ snapshots: '@walletconnect/logger': 2.1.2 '@walletconnect/universal-provider': 2.21.0(bufferutil@4.0.9)(encoding@0.1.13)(ioredis@5.6.1)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76) valtio: 1.13.2(react@18.3.1) - viem: 2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.24.2) + viem: 2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76) transitivePeerDependencies: - '@azure/app-configuration' - '@azure/cosmos' @@ -10072,7 +10072,7 @@ snapshots: '@walletconnect/universal-provider': 2.21.0(bufferutil@4.0.9)(encoding@0.1.13)(ioredis@5.6.1)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76) bs58: 6.0.0 valtio: 1.13.2(react@18.3.1) - viem: 2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.24.2) + viem: 2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76) transitivePeerDependencies: - '@azure/app-configuration' - '@azure/cosmos' @@ -10113,7 +10113,7 @@ snapshots: '@safe-global/safe-apps-sdk@9.1.0(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76)': dependencies: '@safe-global/safe-gateway-typescript-sdk': 3.23.1 - viem: 2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.24.2) + viem: 2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76) transitivePeerDependencies: - bufferutil - typescript @@ -11072,19 +11072,19 @@ snapshots: dependencies: '@types/yargs-parser': 21.0.3 - '@wagmi/connectors@5.11.2(@tanstack/react-query@5.84.1(react@18.3.1))(@wagmi/core@2.21.2(@tanstack/query-core@5.83.1)(react@18.3.1)(typescript@5.8.3)(use-sync-external-store@1.4.0(react@18.3.1))(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76)))(bufferutil@4.0.9)(encoding@0.1.13)(ioredis@5.6.1)(react@18.3.1)(typescript@5.8.3)(use-sync-external-store@1.4.0(react@18.3.1))(utf-8-validate@5.0.10)(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76))(wagmi@2.17.5(@tanstack/query-core@5.83.1)(@tanstack/react-query@5.84.1(react@18.3.1))(bufferutil@4.0.9)(encoding@0.1.13)(ioredis@5.6.1)(react@18.3.1)(typescript@5.8.3)(utf-8-validate@5.0.10)(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76))(zod@3.25.76))(zod@3.25.76)': + '@wagmi/connectors@5.11.2(@tanstack/react-query@5.84.1(react@18.3.1))(@wagmi/core@2.21.2(@tanstack/query-core@5.83.1)(react@18.3.1)(typescript@5.8.3)(use-sync-external-store@1.4.0(react@18.3.1))(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.24.2)))(bufferutil@4.0.9)(encoding@0.1.13)(ioredis@5.6.1)(react@18.3.1)(typescript@5.8.3)(use-sync-external-store@1.4.0(react@18.3.1))(utf-8-validate@5.0.10)(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76))(wagmi@2.17.5(@tanstack/query-core@5.83.1)(@tanstack/react-query@5.84.1(react@18.3.1))(bufferutil@4.0.9)(encoding@0.1.13)(ioredis@5.6.1)(react@18.3.1)(typescript@5.8.3)(utf-8-validate@5.0.10)(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76))(zod@3.25.76))(zod@3.25.76)': dependencies: '@base-org/account': 1.1.1(bufferutil@4.0.9)(react@18.3.1)(typescript@5.8.3)(use-sync-external-store@1.4.0(react@18.3.1))(utf-8-validate@5.0.10)(zod@3.25.76) '@coinbase/wallet-sdk': 4.3.6(bufferutil@4.0.9)(react@18.3.1)(typescript@5.8.3)(use-sync-external-store@1.4.0(react@18.3.1))(utf-8-validate@5.0.10)(zod@3.25.76) - '@gemini-wallet/core': 0.2.0(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76)) + '@gemini-wallet/core': 0.2.0(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.24.2)) '@metamask/sdk': 0.33.1(bufferutil@4.0.9)(encoding@0.1.13)(utf-8-validate@5.0.10) '@safe-global/safe-apps-provider': 0.18.6(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76) '@safe-global/safe-apps-sdk': 9.1.0(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76) - '@wagmi/core': 2.21.2(@tanstack/query-core@5.83.1)(react@18.3.1)(typescript@5.8.3)(use-sync-external-store@1.4.0(react@18.3.1))(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76)) + '@wagmi/core': 2.21.2(@tanstack/query-core@5.83.1)(react@18.3.1)(typescript@5.8.3)(use-sync-external-store@1.4.0(react@18.3.1))(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.24.2)) '@walletconnect/ethereum-provider': 2.21.1(bufferutil@4.0.9)(encoding@0.1.13)(ioredis@5.6.1)(react@18.3.1)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76) cbw-sdk: '@coinbase/wallet-sdk@3.9.3' - porto: 0.2.19(@tanstack/react-query@5.84.1(react@18.3.1))(@wagmi/core@2.21.2(@tanstack/query-core@5.83.1)(react@18.3.1)(typescript@5.8.3)(use-sync-external-store@1.4.0(react@18.3.1))(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76)))(react@18.3.1)(typescript@5.8.3)(use-sync-external-store@1.4.0(react@18.3.1))(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76))(wagmi@2.17.5(@tanstack/query-core@5.83.1)(@tanstack/react-query@5.84.1(react@18.3.1))(bufferutil@4.0.9)(encoding@0.1.13)(ioredis@5.6.1)(react@18.3.1)(typescript@5.8.3)(utf-8-validate@5.0.10)(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76))(zod@3.25.76)) - viem: 2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.24.2) + porto: 0.2.19(@tanstack/react-query@5.84.1(react@18.3.1))(@wagmi/core@2.21.2(@tanstack/query-core@5.83.1)(react@18.3.1)(typescript@5.8.3)(use-sync-external-store@1.4.0(react@18.3.1))(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.24.2)))(react@18.3.1)(typescript@5.8.3)(use-sync-external-store@1.4.0(react@18.3.1))(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.24.2))(wagmi@2.17.5(@tanstack/query-core@5.83.1)(@tanstack/react-query@5.84.1(react@18.3.1))(bufferutil@4.0.9)(encoding@0.1.13)(ioredis@5.6.1)(react@18.3.1)(typescript@5.8.3)(utf-8-validate@5.0.10)(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76))(zod@3.25.76)) + viem: 2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76) optionalDependencies: typescript: 5.8.3 transitivePeerDependencies: @@ -11118,11 +11118,11 @@ snapshots: - wagmi - zod - '@wagmi/core@2.21.2(@tanstack/query-core@5.83.1)(react@18.3.1)(typescript@5.8.3)(use-sync-external-store@1.4.0(react@18.3.1))(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76))': + '@wagmi/core@2.21.2(@tanstack/query-core@5.83.1)(react@18.3.1)(typescript@5.8.3)(use-sync-external-store@1.4.0(react@18.3.1))(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.24.2))': dependencies: eventemitter3: 5.0.1 mipd: 0.0.7(typescript@5.8.3) - viem: 2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.24.2) + viem: 2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76) zustand: 5.0.0(react@18.3.1)(use-sync-external-store@1.4.0(react@18.3.1)) optionalDependencies: '@tanstack/query-core': 5.83.1 @@ -11670,11 +11670,6 @@ snapshots: typescript: 5.8.3 zod: 3.22.4 - abitype@1.0.8(typescript@5.8.3)(zod@3.24.2): - optionalDependencies: - typescript: 5.8.3 - zod: 3.24.2 - abitype@1.0.8(typescript@5.8.3)(zod@3.25.76): optionalDependencies: typescript: 5.8.3 @@ -11685,11 +11680,6 @@ snapshots: typescript: 5.8.3 zod: 3.22.4 - abitype@1.1.1(typescript@5.8.3)(zod@3.24.2): - optionalDependencies: - typescript: 5.8.3 - zod: 3.24.2 - abitype@1.1.1(typescript@5.8.3)(zod@3.25.76): optionalDependencies: typescript: 5.8.3 @@ -14161,21 +14151,6 @@ snapshots: transitivePeerDependencies: - zod - ox@0.8.6(typescript@5.8.3)(zod@3.24.2): - dependencies: - '@adraffy/ens-normalize': 1.11.1 - '@noble/ciphers': 1.3.0 - '@noble/curves': 1.9.6 - '@noble/hashes': 1.8.0 - '@scure/bip32': 1.7.0 - '@scure/bip39': 1.6.0 - abitype: 1.1.1(typescript@5.8.3)(zod@3.24.2) - eventemitter3: 5.0.1 - optionalDependencies: - typescript: 5.8.3 - transitivePeerDependencies: - - zod - ox@0.8.6(typescript@5.8.3)(zod@3.25.76): dependencies: '@adraffy/ens-normalize': 1.11.1 @@ -14415,21 +14390,21 @@ snapshots: pony-cause@2.1.11: {} - porto@0.2.19(@tanstack/react-query@5.84.1(react@18.3.1))(@wagmi/core@2.21.2(@tanstack/query-core@5.83.1)(react@18.3.1)(typescript@5.8.3)(use-sync-external-store@1.4.0(react@18.3.1))(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76)))(react@18.3.1)(typescript@5.8.3)(use-sync-external-store@1.4.0(react@18.3.1))(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76))(wagmi@2.17.5(@tanstack/query-core@5.83.1)(@tanstack/react-query@5.84.1(react@18.3.1))(bufferutil@4.0.9)(encoding@0.1.13)(ioredis@5.6.1)(react@18.3.1)(typescript@5.8.3)(utf-8-validate@5.0.10)(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76))(zod@3.25.76)): + porto@0.2.19(@tanstack/react-query@5.84.1(react@18.3.1))(@wagmi/core@2.21.2(@tanstack/query-core@5.83.1)(react@18.3.1)(typescript@5.8.3)(use-sync-external-store@1.4.0(react@18.3.1))(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.24.2)))(react@18.3.1)(typescript@5.8.3)(use-sync-external-store@1.4.0(react@18.3.1))(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.24.2))(wagmi@2.17.5(@tanstack/query-core@5.83.1)(@tanstack/react-query@5.84.1(react@18.3.1))(bufferutil@4.0.9)(encoding@0.1.13)(ioredis@5.6.1)(react@18.3.1)(typescript@5.8.3)(utf-8-validate@5.0.10)(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76))(zod@3.25.76)): dependencies: - '@wagmi/core': 2.21.2(@tanstack/query-core@5.83.1)(react@18.3.1)(typescript@5.8.3)(use-sync-external-store@1.4.0(react@18.3.1))(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76)) + '@wagmi/core': 2.21.2(@tanstack/query-core@5.83.1)(react@18.3.1)(typescript@5.8.3)(use-sync-external-store@1.4.0(react@18.3.1))(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.24.2)) hono: 4.9.10 idb-keyval: 6.2.2 mipd: 0.0.7(typescript@5.8.3) ox: 0.9.8(typescript@5.8.3)(zod@4.1.11) - viem: 2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.24.2) + viem: 2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76) zod: 4.1.11 zustand: 5.0.8(react@18.3.1)(use-sync-external-store@1.4.0(react@18.3.1)) optionalDependencies: '@tanstack/react-query': 5.84.1(react@18.3.1) react: 18.3.1 typescript: 5.8.3 - wagmi: 2.17.5(@tanstack/query-core@5.83.1)(@tanstack/react-query@5.84.1(react@18.3.1))(bufferutil@4.0.9)(encoding@0.1.13)(ioredis@5.6.1)(react@18.3.1)(typescript@5.8.3)(utf-8-validate@5.0.10)(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76))(zod@3.24.2) + wagmi: 2.17.5(@tanstack/query-core@5.83.1)(@tanstack/react-query@5.84.1(react@18.3.1))(bufferutil@4.0.9)(encoding@0.1.13)(ioredis@5.6.1)(react@18.3.1)(typescript@5.8.3)(utf-8-validate@5.0.10)(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.24.2))(zod@3.24.2) transitivePeerDependencies: - '@types/react' - immer @@ -15301,7 +15276,7 @@ snapshots: utils-merge@1.0.1: {} - uuid@10.0.0: {} + uuid@13.0.0: {} uuid@8.3.2: {} @@ -15359,23 +15334,6 @@ snapshots: - utf-8-validate - zod - viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.24.2): - dependencies: - '@noble/curves': 1.9.2 - '@noble/hashes': 1.8.0 - '@scure/bip32': 1.7.0 - '@scure/bip39': 1.6.0 - abitype: 1.0.8(typescript@5.8.3)(zod@3.24.2) - isows: 1.0.7(ws@8.18.2(bufferutil@4.0.9)(utf-8-validate@5.0.10)) - ox: 0.8.6(typescript@5.8.3)(zod@3.24.2) - ws: 8.18.2(bufferutil@4.0.9)(utf-8-validate@5.0.10) - optionalDependencies: - typescript: 5.8.3 - transitivePeerDependencies: - - bufferutil - - utf-8-validate - - zod - viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76): dependencies: '@noble/curves': 1.9.2 @@ -15397,14 +15355,14 @@ snapshots: dependencies: xml-name-validator: 5.0.0 - wagmi@2.17.5(@tanstack/query-core@5.83.1)(@tanstack/react-query@5.84.1(react@18.3.1))(bufferutil@4.0.9)(encoding@0.1.13)(ioredis@5.6.1)(react@18.3.1)(typescript@5.8.3)(utf-8-validate@5.0.10)(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76))(zod@3.24.2): + wagmi@2.17.5(@tanstack/query-core@5.83.1)(@tanstack/react-query@5.84.1(react@18.3.1))(bufferutil@4.0.9)(encoding@0.1.13)(ioredis@5.6.1)(react@18.3.1)(typescript@5.8.3)(utf-8-validate@5.0.10)(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.24.2))(zod@3.24.2): dependencies: '@tanstack/react-query': 5.84.1(react@18.3.1) - '@wagmi/connectors': 5.11.2(@tanstack/react-query@5.84.1(react@18.3.1))(@wagmi/core@2.21.2(@tanstack/query-core@5.83.1)(react@18.3.1)(typescript@5.8.3)(use-sync-external-store@1.4.0(react@18.3.1))(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76)))(bufferutil@4.0.9)(encoding@0.1.13)(ioredis@5.6.1)(react@18.3.1)(typescript@5.8.3)(use-sync-external-store@1.4.0(react@18.3.1))(utf-8-validate@5.0.10)(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76))(wagmi@2.17.5(@tanstack/query-core@5.83.1)(@tanstack/react-query@5.84.1(react@18.3.1))(bufferutil@4.0.9)(encoding@0.1.13)(ioredis@5.6.1)(react@18.3.1)(typescript@5.8.3)(utf-8-validate@5.0.10)(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76))(zod@3.25.76))(zod@3.25.76) - '@wagmi/core': 2.21.2(@tanstack/query-core@5.83.1)(react@18.3.1)(typescript@5.8.3)(use-sync-external-store@1.4.0(react@18.3.1))(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76)) + '@wagmi/connectors': 5.11.2(@tanstack/react-query@5.84.1(react@18.3.1))(@wagmi/core@2.21.2(@tanstack/query-core@5.83.1)(react@18.3.1)(typescript@5.8.3)(use-sync-external-store@1.4.0(react@18.3.1))(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.24.2)))(bufferutil@4.0.9)(encoding@0.1.13)(ioredis@5.6.1)(react@18.3.1)(typescript@5.8.3)(use-sync-external-store@1.4.0(react@18.3.1))(utf-8-validate@5.0.10)(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76))(wagmi@2.17.5(@tanstack/query-core@5.83.1)(@tanstack/react-query@5.84.1(react@18.3.1))(bufferutil@4.0.9)(encoding@0.1.13)(ioredis@5.6.1)(react@18.3.1)(typescript@5.8.3)(utf-8-validate@5.0.10)(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76))(zod@3.25.76))(zod@3.25.76) + '@wagmi/core': 2.21.2(@tanstack/query-core@5.83.1)(react@18.3.1)(typescript@5.8.3)(use-sync-external-store@1.4.0(react@18.3.1))(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.24.2)) react: 18.3.1 use-sync-external-store: 1.4.0(react@18.3.1) - viem: 2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.24.2) + viem: 2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76) optionalDependencies: typescript: 5.8.3 transitivePeerDependencies: @@ -15611,8 +15569,8 @@ snapshots: '@solana-program/token-2022': 0.4.2(@solana/kit@2.3.0(fastestsmallesttextencoderdecoder@1.0.22)(typescript@5.8.3)(ws@8.18.0(bufferutil@4.0.9)(utf-8-validate@5.0.10)))(@solana/sysvars@2.3.0(fastestsmallesttextencoderdecoder@1.0.22)(typescript@5.8.3)) '@solana/kit': 2.3.0(fastestsmallesttextencoderdecoder@1.0.22)(typescript@5.8.3)(ws@8.18.0(bufferutil@4.0.9)(utf-8-validate@5.0.10)) '@solana/transaction-confirmation': 2.3.0(fastestsmallesttextencoderdecoder@1.0.22)(typescript@5.8.3)(ws@8.18.0(bufferutil@4.0.9)(utf-8-validate@5.0.10)) - viem: 2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.24.2) - wagmi: 2.17.5(@tanstack/query-core@5.83.1)(@tanstack/react-query@5.84.1(react@18.3.1))(bufferutil@4.0.9)(encoding@0.1.13)(ioredis@5.6.1)(react@18.3.1)(typescript@5.8.3)(utf-8-validate@5.0.10)(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76))(zod@3.24.2) + viem: 2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.25.76) + wagmi: 2.17.5(@tanstack/query-core@5.83.1)(@tanstack/react-query@5.84.1(react@18.3.1))(bufferutil@4.0.9)(encoding@0.1.13)(ioredis@5.6.1)(react@18.3.1)(typescript@5.8.3)(utf-8-validate@5.0.10)(viem@2.33.2(bufferutil@4.0.9)(typescript@5.8.3)(utf-8-validate@5.0.10)(zod@3.24.2))(zod@3.24.2) zod: 3.25.76 transitivePeerDependencies: - '@azure/app-configuration' diff --git a/apps/api/src/lib/concurrency-limit.ts b/apps/api/src/lib/concurrency-limit.ts index f09b34c1a..913e4492b 100644 --- a/apps/api/src/lib/concurrency-limit.ts +++ b/apps/api/src/lib/concurrency-limit.ts @@ -340,14 +340,15 @@ export async function concurrentJobDone(job: NuQJob) { abTestJob(nextJob.job.data); - await scrapeQueue.addJob( + await scrapeQueue.promoteJobFromBacklogOrAdd( nextJob.job.id, + nextJob.job.data, { - ...nextJob.job.data, - concurrencyLimitHit: true, + priority: nextJob.job.priority, + listenable: nextJob.job.listenable, + ownerId: nextJob.job.data.team_id ?? undefined, + groupId: nextJob.job.data.crawl_id ?? undefined, }, - nextJob.job.priority, - nextJob.job.listenable, ); } } diff --git a/apps/api/src/services/queue-jobs.ts b/apps/api/src/services/queue-jobs.ts index 885be5005..36a1722ed 100644 --- a/apps/api/src/services/queue-jobs.ts +++ b/apps/api/src/services/queue-jobs.ts @@ -43,6 +43,21 @@ async function _addScrapeJobToConcurrencyQueue( priority: number = 0, listenable: boolean = false, ) { + await scrapeQueue.addJob( + jobId, + { + ...webScraperOptions, + concurrencyLimited: true, + }, + { + priority, + listenable, + ownerId: webScraperOptions.team_id ?? undefined, + groupId: webScraperOptions.crawl_id ?? undefined, + backlogged: true, + }, + ); + await pushConcurrencyLimitedJob( webScraperOptions.team_id, { @@ -57,6 +72,44 @@ async function _addScrapeJobToConcurrencyQueue( ); } +async function _addScrapeJobsToConcurrencyQueue( + jobs: { + data: any; + jobId: string; + priority: number; + listenable?: boolean; + }[], +) { + await scrapeQueue.addJobs( + jobs.map(job => ({ + id: job.jobId, + data: job.data, + options: { + priority: job.priority, + listenable: job.listenable ?? false, + ownerId: job.data.team_id ?? undefined, + groupId: job.data.crawl_id ?? undefined, + backlogged: true, + }, + })), + ); + + for (const job of jobs) { + await pushConcurrencyLimitedJob( + job.data.team_id, + { + id: job.jobId, + data: job.data, + priority: job.priority, + listenable: job.listenable ?? false, + }, + job.data.crawl_id + ? Infinity + : (job.data.scrapeOptions?.timeout ?? 60 * 1000), + ); + } +} + export async function _addScrapeJobToBullMQ( webScraperOptions: ScrapeJobData, jobId: string, @@ -86,7 +139,59 @@ export async function _addScrapeJobToBullMQ( } } - return await scrapeQueue.addJob(jobId, webScraperOptions, priority, listenable); + return await scrapeQueue.addJob(jobId, webScraperOptions, { + priority, + listenable, + ownerId: webScraperOptions.team_id ?? undefined, + groupId: webScraperOptions.crawl_id ?? undefined, + }); +} + +async function _addScrapeJobsToBullMQ( + jobs: { + data: any; + jobId: string; + priority: number; + listenable?: boolean; + }[], +): Promise[]> { + for (const job of jobs) { + if (job.data.mode === "single_urls") { + abTestJob(job.data); + } + + if (job.data && job.data.team_id) { + await pushConcurrencyLimitActiveJob( + job.data.team_id, + job.jobId, + 60 * 1000, + ); // 60s default timeout + + if (job.data.crawl_id) { + const sc = await getCrawl(job.data.crawl_id); + if (sc?.crawlerOptions?.delay || sc?.maxConcurrency) { + await pushCrawlConcurrencyLimitActiveJob( + job.data.crawl_id, + job.jobId, + 60 * 1000, + ); + } + } + } + } + + return await scrapeQueue.addJobs( + jobs.map(job => ({ + id: job.jobId, + data: job.data, + options: { + priority: job.priority, + listenable: job.listenable ?? false, + ownerId: job.data.team_id ?? undefined, + groupId: job.data.crawl_id ?? undefined, + }, + })), + ); } async function addScrapeJobRaw( @@ -181,10 +286,20 @@ async function addScrapeJobRaw( webScraperOptions.concurrencyLimited = true; - await _addScrapeJobToConcurrencyQueue(webScraperOptions, jobId, priority, listenable); + await _addScrapeJobToConcurrencyQueue( + webScraperOptions, + jobId, + priority, + listenable, + ); return null; } else { - return await _addScrapeJobToBullMQ(webScraperOptions, jobId, priority, listenable); + return await _addScrapeJobToBullMQ( + webScraperOptions, + jobId, + priority, + listenable, + ); } } @@ -368,27 +483,22 @@ export async function addScrapeJobs( } } - await Promise.all( - addToCQ.map(async job => { - const size = JSON.stringify(job.data).length; - await _addScrapeJobToConcurrencyQueue( - { ...job.data, traceContext }, - job.jobId, - job.priority, - job.listenable, - ); - }), + await _addScrapeJobsToConcurrencyQueue( + addToCQ.map(job => ({ + jobId: job.jobId, + data: { ...job.data, traceContext }, + priority: job.priority, + listenable: job.listenable, + })), ); - await Promise.all( - addToBull.map(async job => { - await _addScrapeJobToBullMQ( - { ...job.data, traceContext }, - job.jobId, - job.priority, - job.listenable, - ); - }), + await _addScrapeJobsToBullMQ( + addToBull.map(job => ({ + jobId: job.jobId, + data: { ...job.data, traceContext }, + priority: job.priority, + listenable: job.listenable, + })), ); } } diff --git a/apps/api/src/services/worker/nuq.ts b/apps/api/src/services/worker/nuq.ts index 13a49a555..9c857364c 100644 --- a/apps/api/src/services/worker/nuq.ts +++ b/apps/api/src/services/worker/nuq.ts @@ -4,6 +4,7 @@ import { Client, Pool } from "pg"; import { type ScrapeJobData } from "../../types"; import { withSpan, setSpanAttributes } from "../../lib/otel-tracer"; import amqp from "amqplib"; +import { v5 as uuidv5, validate as isUUID } from "uuid"; // === Basics @@ -16,7 +17,12 @@ nuqPool.on("error", err => logger.error("Error in NuQ idle client", { err, module: "nuq" }), ); -export type NuQJobStatus = "queued" | "active" | "completed" | "failed"; // must match nuq.job_status enum +export type NuQJobStatus = + | "queued" + | "active" + | "completed" + | "failed" + | "backlog"; export type NuQJob = { id: string; status: NuQJobStatus; @@ -28,14 +34,39 @@ export type NuQJob = { returnvalue?: ReturnValue; failedReason?: string; lock?: string; + ownerId?: string; + groupId?: string; }; +type NuQJobOptions = { + priority?: number; + listenable?: boolean; + ownerId?: string; + groupId?: string; + backlogged?: boolean; +}; + +type NuQOptions = { + backlog?: boolean; +}; + +// owner IDs can sometimes be non-UUID, so let's normalize it to avoid query breakage - mogery +const normalizedUUIDNamespace = "0f38e00e-d7ee-4b77-8a7a-a787a3537ca2"; +function normalizeOwnerId(ownerId: string | undefined | null): string | null { + if (typeof ownerId !== "string") return null; + if (isUUID(ownerId)) return ownerId; + return uuidv5(ownerId, normalizedUUIDNamespace); +} + const listenChannelId = process.env.NUQ_POD_NAME ?? "main"; // === Queue class NuQ { - constructor(public readonly queueName: string) {} + constructor( + public readonly queueName: string, + public readonly options: NuQOptions, + ) {} // === Listener @@ -85,7 +116,7 @@ class NuQ { let reconnectTimeout: NodeJS.Timeout | null = null; - const onClose = (function onClose() { + const onClose = function onClose() { logger.info("NuQ listener channel closed", { module: "nuq/rabbitmq", }); @@ -104,7 +135,7 @@ class NuQ { 250, ); return; - }).bind(this); + }.bind(this); connection.on("close", onClose); channel.on("close", onClose); @@ -114,7 +145,7 @@ class NuQ { (msg => { if (msg === null) { onClose(); - return; + return; } logger.info("NuQ job received", { @@ -266,7 +297,7 @@ class NuQ { await this.startSender(); if (this.sender) { - await this.sender.channel.sendToQueue( + this.sender.channel.sendToQueue( this.queueName + ".listen." + listenChannelId, Buffer.from(status, "utf8"), { @@ -286,7 +317,7 @@ class NuQ { await this.startSender(); if (this.sender) { - await this.sender.channel.sendToQueue( + this.sender.channel.sendToQueue( this.queueName + ".prefetch", Buffer.from(JSON.stringify(job), "utf8"), { @@ -314,13 +345,28 @@ class NuQ { "returnvalue", "failedreason", "lock", + "owner_id", + "group_id", ]; - private rowToJob(row: any): NuQJob | null { + private readonly jobBacklogReturning = [ + "id", + "created_at", + "priority", + "data", + "listen_channel_id", + "owner_id", + "group_id", + ]; + + private rowToJob( + row: any, + backlogged?: boolean, + ): NuQJob | null { if (!row) return null; return { id: row.id, - status: row.status, + status: backlogged ? "backlog" : row.status, createdAt: new Date(row.created_at), priority: row.priority, data: row.data, @@ -329,6 +375,8 @@ class NuQ { returnvalue: row.returnvalue ?? undefined, failedReason: row.failedreason ?? undefined, lock: row.lock ?? undefined, + ownerId: row.owner_id ?? undefined, + groupId: row.group_id ?? undefined, }; } @@ -503,16 +551,15 @@ class NuQ { public async addJob( id: string, data: JobData, - priority: number = 0, - listenable: boolean = false, + options: NuQJobOptions, ): Promise> { return withSpan("nuq.addJob", async span => { setSpanAttributes(span, { "nuq.queue_name": this.queueName, "nuq.job_id": id, - "nuq.priority": priority, + "nuq.priority": options.priority ?? 0, "nuq.zero_data_retention": (data as any)?.zeroDataRetention ?? false, - "nuq.listenable": listenable, + "nuq.listenable": options.listenable ?? false, }); const start = Date.now(); @@ -520,8 +567,15 @@ class NuQ { const result = this.rowToJob( ( await nuqPool.query( - `INSERT INTO ${this.queueName} (id, data, priority, listen_channel_id) VALUES ($1, $2, $3, $4) RETURNING ${this.jobReturning.join(", ")};`, - [id, data, priority, listenable ? listenChannelId : null], + `INSERT INTO ${this.queueName}${options.backlogged ? "_backlog" : ""} (id, data, priority, listen_channel_id, owner_id, group_id) VALUES ($1, $2, $3, $4, $5, $6) RETURNING ${(options.backlogged ? this.jobBacklogReturning : this.jobReturning).join(", ")};`, + [ + id, + data, + options.priority ?? 0, + options.listenable ? listenChannelId : null, + normalizeOwnerId(options.ownerId), + options.groupId ?? null, + ], ) ).rows[0], )!; @@ -547,6 +601,189 @@ class NuQ { }); } + public async addJobs( + jobs: Array<{ + id: string; + data: JobData; + options: NuQJobOptions; + }>, + ): Promise[]> { + return withSpan("nuq.addJobs", async span => { + setSpanAttributes(span, { + "nuq.queue_name": this.queueName, + "nuq.jobs_count": jobs.length, + }); + + if (jobs.length === 0) { + return []; + } + + const start = Date.now(); + try { + // Separate jobs into backlogged and non-backlogged groups + const regularJobs: typeof jobs = []; + const backloggedJobs: typeof jobs = []; + + for (const job of jobs) { + if (job.options.backlogged) { + backloggedJobs.push(job); + } else { + regularJobs.push(job); + } + } + + const results: NuQJob[] = []; + + // Batch size: 6 params per job, stay well under PG's 65535 param limit + // 1000 jobs = 6000 params, leaving plenty of headroom + const BATCH_SIZE = 1000; + + // Helper function to build and execute bulk insert with batching + const bulkInsert = async ( + jobsToInsert: typeof jobs, + tableSuffix: string, + ) => { + if (jobsToInsert.length === 0) return; + + // Process in batches + for ( + let offset = 0; + offset < jobsToInsert.length; + offset += BATCH_SIZE + ) { + const batch = jobsToInsert.slice(offset, offset + BATCH_SIZE); + + // Build the VALUES clause and parameters array + const valuesPlaceholders: string[] = []; + const params: any[] = []; + + for (let i = 0; i < batch.length; i++) { + const job = batch[i]; + const baseIdx = i * 6 + 1; + + valuesPlaceholders.push( + `($${baseIdx}, $${baseIdx + 1}, $${baseIdx + 2}, $${baseIdx + 3}, $${baseIdx + 4}, $${baseIdx + 5})`, + ); + + params.push( + job.id, + job.data, + job.options.priority ?? 0, + job.options.listenable ? listenChannelId : null, + normalizeOwnerId(job.options.ownerId), + job.options.groupId ?? null, + ); + } + + const query = `INSERT INTO ${this.queueName}${tableSuffix} (id, data, priority, listen_channel_id, owner_id, group_id) VALUES ${valuesPlaceholders.join(", ")} RETURNING ${(tableSuffix === "_backlog" ? this.jobBacklogReturning : this.jobReturning).join(", ")};`; + + const result = await nuqPool.query(query, params); + + // Convert rows to jobs and maintain order + const jobMap = new Map( + result.rows.map(row => [ + row.id, + this.rowToJob(row, tableSuffix === "_backlog")!, + ]), + ); + + for (const job of batch) { + const insertedJob = jobMap.get(job.id); + if (insertedJob) { + results.push(insertedJob); + } + } + } + }; + + // Insert regular jobs + await bulkInsert(regularJobs, ""); + + // Insert backlogged jobs + await bulkInsert(backloggedJobs, "_backlog"); + + setSpanAttributes(span, { + "nuq.jobs_created": results.length, + "nuq.regular_jobs_count": regularJobs.length, + "nuq.backlogged_jobs_count": backloggedJobs.length, + }); + + return results; + } finally { + const duration = Date.now() - start; + setSpanAttributes(span, { + "nuq.duration_ms": duration, + }); + logger.info("nuqAddJobs metrics", { + module: "nuq/metrics", + method: "nuqAddJobs", + duration, + jobsCount: jobs.length, + }); + } + }); + } + + public async promoteJobFromBacklogOrAdd( + id: string, + data: JobData, + options: NuQJobOptions, + ): Promise> { + return withSpan("nuq.promoteJobFromBacklogOrAdd", async span => { + setSpanAttributes(span, { + "nuq.queue_name": this.queueName, + "nuq.job_id": id, + "nuq.priority": options.priority ?? 0, + "nuq.zero_data_retention": (data as any)?.zeroDataRetention ?? false, + "nuq.listenable": options.listenable ?? false, + }); + + const start = Date.now(); + try { + const result = this.rowToJob( + ( + await nuqPool.query( + ` + WITH ins AS ( + INSERT INTO ${this.queueName} (id, data, created_at, priority, listen_channel_id, owner_id, group_id) + SELECT b.id, b.data, b.created_at, b.priority, b.listen_channel_id, b.owner_id, b.group_id + FROM ${this.queueName}_backlog b + WHERE b.id = $1 + LIMIT 1 + RETURNING ${this.jobReturning.join(", ")} + ), del AS ( + DELETE FROM ${this.queueName}_backlog + WHERE id = $1 + ) + SELECT * FROM ins + `, + [id], + ) + ).rows[0], + ); + + if (!result) { + return await this.addJob(id, data, { + ...options, + backlogged: false, + }); + } + + return result; + } finally { + const duration = Date.now() - start; + setSpanAttributes(span, { + "nuq.duration_ms": duration, + }); + logger.info("nuqPromoteJobFromBacklogOrAdd metrics", { + module: "nuq/metrics", + method: "nuqPromoteJobFromBacklogOrAdd", + duration, + }); + } + }); + } + private readonly nuqWaitMode = process.env.NUQ_WAIT_MODE === "listen" || process.env.NUQ_RABBITMQ_URL ? ("listen" as const) @@ -944,7 +1181,9 @@ export async function nuqHealthCheck(): Promise { // === Instances -export const scrapeQueue = new NuQ("nuq.queue_scrape"); +export const scrapeQueue = new NuQ("nuq.queue_scrape", { + backlog: true, +}); // === Cleanup diff --git a/apps/api/src/services/worker/scrape-worker.ts b/apps/api/src/services/worker/scrape-worker.ts index 14b888a40..0d4804b82 100644 --- a/apps/api/src/services/worker/scrape-worker.ts +++ b/apps/api/src/services/worker/scrape-worker.ts @@ -37,7 +37,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 { createWebhookSender, WebhookEvent } from "../webhook"; +import { createWebhookSender, WebhookEvent } from "../webhook/index"; import { CustomError } from "../../lib/custom-error"; import { startWebScraperPipeline } from "../../main/runWebScraper"; import { CostTracking } from "../../lib/cost-tracking"; @@ -338,10 +338,7 @@ async function processJob(job: NuQJob) { // Store robots blocked URLs in Redis set for (const [url, reason] of links.denialReasons) { if (reason === "URL blocked by robots.txt") { - await recordRobotsBlocked( - job.data.crawl_id, - url - ); + await recordRobotsBlocked(job.data.crawl_id, url); } } @@ -552,15 +549,12 @@ async function processJob(job: NuQJob) { error instanceof Error && error.message === "URL blocked by robots.txt" ) { - await recordRobotsBlocked( - job.data.crawl_id, - job.data.url, - ); + await recordRobotsBlocked(job.data.crawl_id, job.data.url); } } catch (e) { logger.debug("Failed to record top-level robots block", { e }); } - + if (job.data.crawl_id) { const sc = (await getCrawl(job.data.crawl_id)) as StoredCrawl; @@ -1182,7 +1176,7 @@ async function processJobWithTracing(job: NuQJob, logger: any) { await concurrentJobDone(job); } } catch (error) { - logger.debug("Job failed", { error }); + logger.warn("Job failed", { error }); Sentry.captureException(error); if (error instanceof TransportableError) { throw new Error(serializeTransportableError(error)); diff --git a/apps/nuq-postgres/nuq.sql b/apps/nuq-postgres/nuq.sql index f22a6b83a..df49dad3d 100644 --- a/apps/nuq-postgres/nuq.sql +++ b/apps/nuq-postgres/nuq.sql @@ -22,6 +22,8 @@ CREATE TABLE IF NOT EXISTS nuq.queue_scrape ( listen_channel_id text, -- for listenable jobs over rabbitmq returnvalue jsonb, -- only for selfhost failedreason text, -- only for selfhost + owner_id uuid, + group_id uuid, CONSTRAINT queue_scrape_pkey PRIMARY KEY (id) ); @@ -36,6 +38,17 @@ CREATE INDEX IF NOT EXISTS nuq_queue_scrape_queued_optimal_2_idx ON nuq.queue_sc CREATE INDEX IF NOT EXISTS nuq_queue_scrape_failed_created_at_idx ON nuq.queue_scrape USING btree (created_at) WHERE (status = 'failed'::nuq.job_status); CREATE INDEX IF NOT EXISTS nuq_queue_scrape_completed_created_at_idx ON nuq.queue_scrape USING btree (created_at) WHERE (status = 'completed'::nuq.job_status); +CREATE TABLE IF NOT EXISTS nuq.queue_scrape_backlog ( + id uuid NOT NULL DEFAULT gen_random_uuid(), + data jsonb, + created_at timestamp with time zone NOT NULL DEFAULT now(), + priority int NOT NULL DEFAULT 0, + listen_channel_id text, -- for listenable jobs over rabbitmq + owner_id uuid, + group_id uuid, + CONSTRAINT queue_scrape_backlog_pkey PRIMARY KEY (id) +); + SELECT cron.schedule('nuq_queue_scrape_clean_completed', '*/5 * * * *', $$ DELETE FROM nuq.queue_scrape WHERE nuq.queue_scrape.status = 'completed'::nuq.job_status AND nuq.queue_scrape.created_at < now() - interval '1 hour'; $$);