rag-service/src/modules/ocr/dispatcher.ts

108 lines
4.3 KiB
TypeScript

import type { OcrJobRow } from "../catalog/repository.js";
import { OcrClientError, type OcrClient, type OcrResult } from "./client.js";
export interface OcrDispatchStore {
claimNextOcrJob(leaseMs: number): Promise<OcrJobRow | undefined>;
claimOcrJob(jobId: string, leaseMs: number): Promise<OcrJobRow | undefined>;
recoverExpiredOcrLeases(): Promise<OcrJobRow[]>;
setOcrRemoteJob(jobId: string, remoteJobId: string, leaseMs: number): Promise<void>;
requeueOcrJob(jobId: string, code: string, detail: string, delayMs: number): Promise<void>;
completeOcrJob(jobId: string, result: OcrResult): Promise<boolean>;
failOcrJob(jobId: string, code: string, detail: string): Promise<void>;
markReviewRequired(versionId: string): Promise<void>;
markFailed(versionId: string, code: string, detail: string): Promise<void>;
}
export type OcrDispatchInput = { bytes: Buffer; documentSha256: string };
export type OcrDispatchResult = "idle" | "pending" | "succeeded" | "failed";
export class OcrDispatcher {
private activeDrain: Promise<number> | undefined;
constructor(
private readonly store: OcrDispatchStore,
private readonly client: OcrClient,
private readonly loadInput: (job: OcrJobRow) => Promise<OcrDispatchInput>,
private readonly leaseMs = 30_000
) {}
async runOnce(): Promise<OcrDispatchResult> {
const job = await this.store.claimNextOcrJob(this.leaseMs);
if (!job) return "idle";
return this.dispatch(job);
}
private async dispatch(job: OcrJobRow): Promise<OcrDispatchResult> {
try {
const input = await this.loadInput(job);
let remoteJobId = job.remoteJobId;
if (!remoteJobId) {
const acknowledgement = await this.client.submit(input.bytes, {
documentSha256: input.documentSha256,
pages: job.requestedPages
}, job.remoteIdempotencyKey);
remoteJobId = acknowledgement.jobId;
await this.store.setOcrRemoteJob(job.jobId, remoteJobId, this.leaseMs);
}
const status = await this.client.getStatus(remoteJobId);
if (status.status === "queued" || status.status === "running") {
const delayMs = Math.min(2_000 * 2 ** Math.max(0, job.attemptCount - 1), 15_000);
await this.store.requeueOcrJob(job.jobId, "OCR_PENDING", status.status, delayMs);
return "pending";
}
if (status.status === "failed") {
const code = status.error?.code ?? "OCR_REMOTE_FAILED";
await this.fail(job, code, status.error?.message ?? "OCR processing failed");
return "failed";
}
const result = await this.client.getResult(remoteJobId, {
documentSha256: input.documentSha256,
pages: job.requestedPages
});
const versionComplete = await this.store.completeOcrJob(job.jobId, result);
if (versionComplete) await this.store.markReviewRequired(job.versionId);
return "succeeded";
} catch (error) {
const detail = error instanceof Error ? error.message : "Unknown OCR dispatch failure";
if (error instanceof OcrClientError && error.retryable) {
await this.store.requeueOcrJob(job.jobId, error.code, detail, 15_000);
return "pending";
}
const code = error instanceof OcrClientError ? error.code : "OCR_DISPATCH_FAILED";
await this.fail(job, code, detail);
return "failed";
}
}
async recoverExpiredLeases(): Promise<number> {
const recovered = await this.store.recoverExpiredOcrLeases();
for (const job of recovered) {
const claimed = await this.store.claimOcrJob(job.jobId, this.leaseMs);
if (claimed) await this.dispatch(claimed);
}
return recovered.length;
}
dispatchAvailable(): Promise<number> {
if (this.activeDrain) return this.activeDrain;
const drain = this.drainAvailable();
this.activeDrain = drain;
const clear = () => { if (this.activeDrain === drain) this.activeDrain = undefined; };
void drain.then(clear, clear);
return drain;
}
private async drainAvailable(): Promise<number> {
let processed = 0;
while (await this.runOnce() !== "idle") processed += 1;
return processed;
}
private async fail(job: OcrJobRow, code: string, detail: string): Promise<void> {
await this.store.failOcrJob(job.jobId, code, detail);
await this.store.markFailed(job.versionId, code, detail);
}
}