Repository navigation
Processor service typing + tech debt (#441, #431, #436) - #506
Conversation
…441, #431, #436) Define per-queue job data and result interfaces (IngestFileJobData, ImageOptimizeJobData, TextExtractJobData, OCRJobData + corresponding result types), replacing all `any` casts in queue-manager workers, addJob/getJob/getQueueStatus methods, worker return types, and db.ts pool/metadata. Remove the stubbed 501 /upload/chunk endpoint and its doc references. Fix the dedupe upload path to return real job IDs from queueProcessingJobs instead of hardcoded { ingest: true } stub. Adds 28 new tests covering type contracts, state mapping, dedupe response shape, and endpoint removal. Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
|
You have reached your Codex usage limits for code reviews. You can see your limits in the Codex usage dashboard. |
|
Warning Rate limit exceeded
⌛ How to resolve this issue?After the wait time has elapsed, a review can be triggered using the We recommend that you space out your commits to avoid hitting the rate limit. 🚦 How do rate limits work?CodeRabbit enforces hourly rate limits for each developer per organization. Our paid plans have higher rate limits than the trial, open-source and free plans. In all cases, we re-allow further reviews after a brief timeout. Please see our FAQ for further information. 📝 WalkthroughWalkthroughRemoves the legacy POST /chunk endpoint and getQueuedJobs helper; refactors upload handlers to return actual queued job IDs; introduces typed job/data/result interfaces and QueueName, updates workers and QueueManager to use those types, tightens DB typings, and adds tests for upload behavior and type contracts. Changes
Sequence Diagram(s)sequenceDiagram
rect rgba(200,220,255,0.5)
participant Client
end
rect rgba(200,255,200,0.5)
participant API as Processor API
end
rect rgba(255,220,200,0.5)
participant QM as QueueManager
participant PgB as PgBoss
end
rect rgba(255,200,255,0.5)
participant Worker
participant DB as Database/ContentStore
end
Client->>API: POST /upload (file)
API->>QM: queueProcessingJobs / addJob(...)
QM->>PgB: send job
PgB-->>QM: jobId
QM-->>API: jobId(s)
API-->>Client: { success: true, jobs: [jobId,...], deduplicated: bool }
PgB->>Worker: deliver job
Worker->>DB: read/write file metadata, content, pages
Worker-->>PgB: complete/failed
PgB-->>QM: monitor-states update
Estimated code review effort🎯 4 (Complex) | ⏱️ ~60 minutes Possibly related issues
Possibly related PRs
Poem
🚥 Pre-merge checks | ✅ 2 | ❌ 1❌ Failed checks (1 warning)
✅ Passed checks (2 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing touches🧪 Generate unit tests (beta)
Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
There was a problem hiding this comment.
Actionable comments posted: 1
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (5)
docs/2.0-architecture/2.2-backend/processor-service.md (2)
103-109:⚠️ Potential issue | 🟡 MinorDocumentation still references
anytypes and outdated API signatures.Lines 106 and 108 still show
data: any,options?: any, andRecord<string, any>in theQueueManagerinterface documentation. Now that the codebase uses strongly-typedaddJob<Q extends QueueName>(queue: Q, data: JobDataMap[Q])andgetQueueStatus(): Promise<Record<QueueName, QueueStats>>, this section is stale. Similarly, line 128 hasmetadata?: anyfor the text extractor.Consider updating these snippets to reflect the new typed contracts so the doc stays in sync with the code.
179-187:⚠️ Potential issue | 🟡 MinorUpload response
jobsshape may be stale.The PR changes the dedupe upload path to return real
string[]job IDs fromqueueProcessingJobs()instead of a hardcoded{ ingest: true }stub. The documented response shape on line 185 (jobs: { ingest?: true } | { textExtraction?: true; ... }) likely no longer matches the actual API response.apps/processor/src/workers/text-extractor.ts (2)
80-100:⚠️ Potential issue | 🟠 MajorThree
as anycasts remain inextractPdfText— contradicts the PR goal of removing allanyfrom workers.Lines 82, 94, and 100 still use
any:
(pdfjsLib as any).getDocument(...)— the pdfjs-dist types may be incomplete, but a narrower cast or a typed wrapper would be cleaner.(item: any) => item.str— pdfjs text content items have a known shape ({ str: string; ... }).metadata.info as any— info can be typed with the known fields used below.Suggested narrowing
+// Narrow type for pdfjs text content items +interface PdfTextItem { + str: string; +} + async function extractPdfText(buffer: Buffer): Promise<{ text: string; metadata: Record<string, unknown> }> { const uint8Array = new Uint8Array(buffer); - const loadingTask = (pdfjsLib as any).getDocument({ data: uint8Array, disableWorker: true }); + const loadingTask = pdfjsLib.getDocument({ data: uint8Array, disableWorker: true }); const pdf = await loadingTask.promise; const metadata = await pdf.getMetadata(); // ... const pageText = textContent.items - .map((item: any) => item.str) + .map((item: PdfTextItem) => item.str) .join(' '); // ... - const info = metadata.info as any; + const info = metadata.info as Record<string, string | undefined>;As per coding guidelines: "Never use
anytypes - always use proper TypeScript types."
56-62: 🛠️ Refactor suggestion | 🟠 MajorUnused variable
textCachePathand inlinerequire()calls.
textCachePathon line 56 is assigned but never used — the actual file path is built fromcacheDirand a hardcoded filename on lines 59-60.Additionally,
require('path')andrequire('fs')on lines 57-61 should be top-level ESM imports per coding guidelines.Suggested fix
Add at the top of the file:
import path from 'path'; import fs from 'fs/promises';Then replace lines 56-62:
- const textCachePath = `${contentHash}/extracted-text.txt`; - const cacheDir = require('path').dirname(await contentStore.getCachePath(contentHash, 'text')); - await require('fs').promises.mkdir(cacheDir, { recursive: true }); - await require('fs').promises.writeFile( - require('path').join(cacheDir, 'extracted-text.txt'), + const cacheDir = path.dirname(await contentStore.getCachePath(contentHash, 'text')); + await fs.mkdir(cacheDir, { recursive: true }); + await fs.writeFile( + path.join(cacheDir, 'extracted-text.txt'), extractedText );As per coding guidelines: "Use ESM modules throughout the codebase."
apps/processor/src/workers/queue-manager.ts (1)
175-198:⚠️ Potential issue | 🟠 Major
addJobreturn type hides a potentialnullfrom PgBoss.
PgBoss.send()returnsPromise<string | null>— it returnsnullwhen a job with the same unique key already exists or when throttled. Theas stringcast on line 197 silently hides this, and callers ofaddJobwill receivenulltyped asstring, violating the function'sPromise<string>signature.Handle the null case explicitly
const jobId = await this.boss.send(queue, data, jobOptions); + if (!jobId) { + throw new Error(`Failed to queue job on ${queue} (duplicate or rejected)`); + } console.log(`Queued job ${jobId} on ${queue}`); - - return jobId as string; + return jobId;
🤖 Fix all issues with AI agents
In `@apps/processor/src/db.ts`:
- Around line 1-5: The PgPoolClient.query result type misdeclares rowCount as
number but pg types it as number | null; update the PgPoolClient interface's
query return type so rowCount is number | null, and then adjust any code using
result.rowCount (e.g., the check in the function that inspects result.rowCount)
to handle null safely—either coerce null to 0 before comparing or explicitly
check for null (use result.rowCount ?? 0 === 0 or result.rowCount !== null &&
result.rowCount === 0) to avoid unexpected behavior.
🧹 Nitpick comments (8)
apps/processor/src/api/upload.ts (1)
415-430:queueProcessingJobsalways creates an array but only ever pushes one element.The function creates a
jobIdsarray, pushes a single ID, and returns it. This is fine if future queues will be added here, but currently the array indirection is unnecessary. Consider a// TODO: add image-optimize / text-extract jobs herecomment to signal intent, or simplify toreturn [jobId].apps/processor/src/api/__tests__/upload-chunk.test.ts (1)
5-17: Source-scanning tests are inherently brittle — consider supplementing with a behavioral test.These tests read
upload.tsas a string and assert on substrings. They'll silently pass if the file is moved/renamed, and they can't catch a re-introduced chunk route registered from a different module. A supertest-based integration test (mounting the router and assertingPOST /chunkreturns 404) would be more robust. That said, as a quick regression guard this is acceptable for now.apps/processor/src/api/__tests__/upload-dedupe.test.ts (1)
15-18: String match is sensitive to whitespace/formatting changes.
'const jobs = await queueProcessingJobs('will break if a formatter rewraps the line or if the variable binding changes tolet. A regex match (e.g.,/const\s+jobs\s*=\s*await\s+queueProcessingJobs\s*\(/) would be slightly more resilient if you want to keep the source-scanning approach.apps/processor/src/types/index.ts (1)
45-49: Consider narrowingpresetto the known preset keys.
presetis typed asstring, butIMAGE_PRESETSdefines a fixed set of keys. A tighter type likekeyof typeof IMAGE_PRESETS(or a dedicated union) would catch invalid presets at compile time rather than runtime. Low priority since the runtime validation already guards this.apps/processor/src/workers/image-processor.ts (1)
6-9: Biome false positive — control characters in regex are intentional here.Biome flags
\x00-\x1fand\x7f-\x9fas unexpected control characters in a regular expression (noControlCharactersInRegex). This is a deliberate sanitization regex stripping C0/C1 control chars and newlines from values before logging to prevent log injection. Consider suppressing the lint rule inline if your CI enforces Biome:Suppress inline
+// biome-ignore lint/suspicious/noControlCharactersInRegex: intentional log-injection sanitization return String(value).replace(/[\x00-\x1f\x7f-\x9f\n\r]/g, '').slice(0, 200);apps/processor/src/workers/__tests__/queue-manager.test.ts (1)
35-53: DuplicatedmapJobStatelogic may drift from the actual implementation.The local
mapJobStatefunction mirrors the privateQueueManager.mapJobStatemethod. If someone updates one without the other, these tests would still pass while the actual behavior changes. Consider exporting the mapping as a standalone utility or a static method so the test can exercise the real code.Example: extract to a shared utility
In
queue-manager.ts:- private mapJobState(state: string): ProcessingJob['status'] { + // Extracted so tests can cover the real implementation + static mapJobState(state: string): ProcessingJob['status'] {In the test:
- const mapJobState = (state: string): 'pending' | 'processing' | 'completed' | 'failed' => { - switch (state) { /* ... */ } - }; + // Import and test the real method + const { QueueManager } = await import('../queue-manager'); + const mapJobState = QueueManager.mapJobState;apps/processor/src/db.ts (2)
12-13:require('pg')violates the ESM-throughout guideline.The codebase guidelines require ESM modules. A dynamic
import()achieves the same lazy-loading without CommonJS:Switch to dynamic ESM import
-// eslint-disable-next-line `@typescript-eslint/no-var-requires` -const { Pool } = require('pg') as { Pool: new (config: { connectionString: string; max: number }) => PgPool }; +const { Pool } = (await import('pg')) as unknown as { Pool: new (config: { connectionString: string; max: number }) => PgPool };Note: if top-level
awaitisn't available, you can move the import insidegetPool():async function getPool(): Promise<PgPool> { if (!pool) { const { Pool } = await import('pg'); // ... } return pool; }This would require making all callers
await getPool()(which they already do via.connect()).As per coding guidelines: "Use ESM modules throughout the codebase".
81-99:getPageForIngestionuses a raw SQL alias — fragile coupling to schema.The query aliases
"filePath"→"contentHash"(line 90). If the underlying column is renamed or the aliasing is wrong, this won't produce a type error — theascast on line 94 silently forces the shape. This is inherent to raw SQL but worth noting as a maintenance risk given that the rest of the monorepo uses Drizzle for schema-aware queries. Based on learnings: "Always use Drizzle client and queries frompagespace/dbfor database access instead of direct queries."
…docs - Fix rowCount type in PgPoolClient (number → number | null), use rows.length check - Replace inline require() calls in text-extractor with top-level ESM imports - Remove unused textCachePath variable in text-extractor - Handle null from PgBoss.send() in addJob instead of casting to string - Replace source-reading tests with behavioral supertest tests (upload-chunk, upload-dedupe) - Extract mapJobState to module-level export, test real function with Riteway - Add pdfjs.ts type interfaces to eliminate 3 any casts in text-extractor - Add Riteway assert helper for processor package - Update processor-service.md docs to match new typed contracts Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
There was a problem hiding this comment.
Actionable comments posted: 1
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (2)
docs/2.0-architecture/2.2-backend/processor-service.md (1)
152-153:⚠️ Potential issue | 🟡 MinorDoc/code mismatch:
pageIdvsfileIdin ingest job data.Line 153 documents the ingest job data as containing
pageId, but inqueue-manager.ts(Line 96) the worker destructuresfileIdfromIngestFileJobData. Align the documentation with the actual type to avoid confusion.apps/processor/src/workers/queue-manager.ts (1)
90-96:⚠️ Potential issue | 🔴 CriticalBatch workers silently mark all jobs as completed even when only the first job is processed — critical data loss.
With
batchSize: 2, pg-boss fetches up to 2 jobs and delivers them as an array to the handler. The code extractsjobs[0]and processes only that job. When the handler resolves (returns without throwing), pg-boss automatically marks all jobs in the batch as completed. This means jobs[1] through jobs[n] are permanently marked completed without being processed, with no retry possible—they are lost.The same issue occurs for:
image-optimize(batchSize: 5, line 163-169): processes only jobs[0]; jobs[1–4] losttext-extract(batchSize: 3, line 173-179): processes only jobs[0]; jobs[1–2] lostThe
ocr-processworker (batchSize: 1, line 183-190) is not affected.Fix: Either set
batchSize: 1for all multi-job workers (simplest), or refactor handlers to process all jobs in the batch with a loop. Current code should not be merged.
🤖 Fix all issues with AI agents
In `@apps/processor/src/workers/queue-manager.ts`:
- Around line 249-254: The counts are wrong because getQueueSize(...) is being
called with the { before: ... } option; switch to querying exact states and
rename variables to match: call this.boss.getQueueSize(queue, { state: 'created'
}) for created, this.boss.getQueueSize(queue, { state: 'active' }) for active,
this.boss.getQueueSize(queue, { state: 'completed' }) for completed and
this.boss.getQueueSize(queue, { state: 'failed' }) for failed (update the
variables in queue-manager.ts accordingly and remove the incorrect { before: ...
} calls to get accurate per-state counts).
🧹 Nitpick comments (2)
apps/processor/src/db.ts (2)
12-13: Prefer ESM import overrequire('pg').The coding guidelines mandate ESM modules throughout the codebase. You can achieve the same type narrowing with an ESM import:
Suggested fix
-// eslint-disable-next-line `@typescript-eslint/no-var-requires` -const { Pool } = require('pg') as { Pool: new (config: { connectionString: string; max: number }) => PgPool }; +import pg from 'pg'; +const Pool = pg.Pool as unknown as new (config: { connectionString: string; max: number }) => PgPool;This eliminates the eslint-disable comment and aligns with the ESM-only convention. As per coding guidelines,
**/*.{ts,tsx}: "Use ESM modules throughout the codebase."
89-99: SQL aliasesfilePathascontentHash— verify this is intentional.The query selects
"filePath"but aliases it as"contentHash", which the return type declares ascontentHash: string. IffilePathis indeed the content hash stored in the DB column, this works but is potentially confusing to future maintainers. A brief inline comment explaining the mapping would help.
Review Feedback SummaryAddressed in 2fbd089:
Addressed in dea9010:
Addressed in 172a1d0:
|
Aligns documentation with actual IngestFileJobData type which uses fileId. Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
There was a problem hiding this comment.
Actionable comments posted: 2
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (2)
apps/processor/src/workers/queue-manager.ts (1)
242-260:⚠️ Potential issue | 🟡 Minor
getQueueStatusis synchronous — good design, but verify docs match.Using the cached monitor-states event is a smart choice to avoid per-call DB queries. Note that the documentation (
processor-service.mdLine 108) declares the return type asPromise<Record<QueueName, QueueStats>>, but the implementation is synchronous. The doc should be updated to drop thePromisewrapper.docs/2.0-architecture/2.2-backend/processor-service.md (1)
293-293:⚠️ Potential issue | 🟡 MinorStale
pageIdreference — should befileId.Line 153 correctly documents
{ fileId, contentHash, mimeType, originalName }, but this line still sayspageId. The PR commit messages explicitly note this was fixed (pageId → fileId), but this occurrence was missed.Proposed fix
-3. **Enqueue**: Queue `ingest-file { pageId, contentHash, mimeType }` +3. **Enqueue**: Queue `ingest-file { fileId, contentHash, mimeType, originalName }`
🤖 Fix all issues with AI agents
In `@apps/processor/src/workers/queue-manager.ts`:
- Line 35: EMPTY_STATS is a single shared QueueStats object that callers can
mutate, corrupting the sentinel across queues; update all places that return or
assign EMPTY_STATS (notably in getQueueStatus and the other usages around the
246-256 region) to return a fresh copy instead of the original reference (e.g.,
replace direct returns/assignments of EMPTY_STATS with a shallow copy like
{...EMPTY_STATS} or provide a factory function newEmptyStats() that constructs a
new QueueStats object) so every consumer receives an independent object.
In `@docs/2.0-architecture/2.2-backend/processor-service.md`:
- Line 108: The docs show getQueueStatus as returning a Promise but the actual
implementation in queue-manager.ts (function getQueueStatus) returns a
synchronous Record<QueueName, QueueStats>; update the documented signature to
match the implementation by changing getQueueStatus(): Promise<Record<QueueName,
QueueStats>> to the synchronous form getQueueStatus(): Record<QueueName,
QueueStats> (alternatively, if you prefer an async API, modify the
getQueueStatus implementation to return Promise.resolve(...) instead—pick one
approach and make doc and the getQueueStatus implementation consistent).
🧹 Nitpick comments (1)
apps/processor/src/workers/queue-manager.ts (1)
220-240:getJobreturnserror: undefinedinstead of omitting the key when no error exists.Line 236:
(output?.error as string) ?? undefinedalways sets theerrorproperty (toundefined), so the returned object will haveerror: undefinedrather than the key being absent. This is a minor semantic difference —ProcessingJob.erroris typed as optional (error?: string), so callers doing'error' in jobwould gettrueeven when there's no error. Not a functional bug, but worth noting if downstream code checks key presence.Proposed: conditionally include `error`
return { id: job.id, type: job.name as QueueName, fileId: (data?.fileId as string) ?? '', contentHash: (data?.contentHash as string) ?? '', status: mapJobState(job.state), result: output as ProcessingJob['result'], - error: (output?.error as string) ?? undefined, + ...(output?.error ? { error: output.error as string } : {}), createdAt: job.createdOn, completedAt: job.completedOn || undefined };
| } | ||
| } | ||
|
|
||
| const EMPTY_STATS: QueueStats = { active: 0, pending: 0, completed: 0, failed: 0 }; |
There was a problem hiding this comment.
Shared EMPTY_STATS reference can be mutated by callers.
EMPTY_STATS is a single object. When multiple queues fall back to it in getQueueStatus, they share the same reference. Any consumer that mutates the returned stats (e.g., status['ingest-file'].active++) will corrupt the shared sentinel for all subsequent calls and all other queues that used the fallback.
Proposed fix — spread a fresh copy
status[queue] = q
? {
pending: q.created + q.retry,
active: q.active,
completed: q.completed,
failed: q.cancelled + q.failed,
}
- : EMPTY_STATS;
+ : { ...EMPTY_STATS };Also applies to: 246-256
🤖 Prompt for AI Agents
In `@apps/processor/src/workers/queue-manager.ts` at line 35, EMPTY_STATS is a
single shared QueueStats object that callers can mutate, corrupting the sentinel
across queues; update all places that return or assign EMPTY_STATS (notably in
getQueueStatus and the other usages around the 246-256 region) to return a fresh
copy instead of the original reference (e.g., replace direct returns/assignments
of EMPTY_STATS with a shallow copy like {...EMPTY_STATS} or provide a factory
function newEmptyStats() that constructs a new QueueStats object) so every
consumer receives an independent object.
| addJob<Q extends QueueName>(queue: Q, data: JobDataMap[Q], options?: PgBoss.SendOptions): Promise<string>; | ||
| getJob(jobId: string): Promise<ProcessingJob | null>; | ||
| getQueueStatus(): Promise<Record<string, any>>; | ||
| getQueueStatus(): Promise<Record<QueueName, QueueStats>>; |
There was a problem hiding this comment.
Doc/code mismatch: getQueueStatus is synchronous, not Promise-returning.
The implementation in queue-manager.ts (Line 242) returns Record<QueueName, QueueStats> synchronously from cached state. This doc shows Promise<Record<QueueName, QueueStats>>.
Proposed fix
- getQueueStatus(): Promise<Record<QueueName, QueueStats>>;
+ getQueueStatus(): Record<QueueName, QueueStats>;📝 Committable suggestion
‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.
| getQueueStatus(): Promise<Record<QueueName, QueueStats>>; | |
| getQueueStatus(): Record<QueueName, QueueStats>; |
🤖 Prompt for AI Agents
In `@docs/2.0-architecture/2.2-backend/processor-service.md` at line 108, The docs
show getQueueStatus as returning a Promise but the actual implementation in
queue-manager.ts (function getQueueStatus) returns a synchronous
Record<QueueName, QueueStats>; update the documented signature to match the
implementation by changing getQueueStatus(): Promise<Record<QueueName,
QueueStats>> to the synchronous form getQueueStatus(): Record<QueueName,
QueueStats> (alternatively, if you prefer an async API, modify the
getQueueStatus implementation to return Promise.resolve(...) instead—pick one
approach and make doc and the getQueueStatus implementation consistent).
- Remove batchSize > 1 from all workers — previous code only processed jobs[0], silently marking remaining batch jobs as completed - Replace incorrect getQueueSize subtraction math with PgBoss monitor-states event cache (already emitted every 30s) - getQueueStatus is now synchronous, serving cached per-queue breakdowns Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
172a1d0 to
58faa5b
Compare
Summary
IngestFileJobData,ImageOptimizeJobData,TextExtractJobData,OCRJobData) and result interfaces (IngestResult,ImageProcessResult,TextExtractResult,OCRResult). GenericaddJob<Q>(), typedgetQueueStatus(), typedgetJob(). Remove allas anycasts from queue-manager workers, worker return types, anddb.tspool/metadata. Also removes deadprocessingTimereference in optimize.ts that was hidden byany.string[]job IDs fromqueueProcessingJobs()instead of hardcoded{ ingest: true }stub. Both upload paths now return identicaljobsshape.Test plan
src/*/__tests__/npx vitest run)src/(npx tsc --noEmit | grep "^src/"→ 0 results)pnpm devand a test upload🤖 Generated with Claude Code
Summary by CodeRabbit
Bug Fixes
API
Types
Documentation
Tests
Chores