feat(ocr): add repository job transitions
This commit is contained in:
parent
936881b38f
commit
d2ecf19102
3 changed files with 314 additions and 1 deletions
|
|
@ -9,7 +9,7 @@
|
||||||
"start": "node dist/server.js",
|
"start": "node dist/server.js",
|
||||||
"migrate:lifecycle": "node dist/modules/catalog/migrations.js",
|
"migrate:lifecycle": "node dist/modules/catalog/migrations.js",
|
||||||
"migrate:legacy:lifecycle": "node dist/scripts/migrate-legacy-lifecycle.js",
|
"migrate:legacy:lifecycle": "node dist/scripts/migrate-legacy-lifecycle.js",
|
||||||
"test": "NODE_ENV=test tsx --test tests/**/*.test.ts",
|
"test": "NODE_ENV=test tsx --test tests/*.test.ts tests/**/*.test.ts",
|
||||||
"check": "tsc --noEmit -p tsconfig.json"
|
"check": "tsc --noEmit -p tsconfig.json"
|
||||||
},
|
},
|
||||||
"dependencies": {
|
"dependencies": {
|
||||||
|
|
|
||||||
|
|
@ -70,6 +70,35 @@ export interface OrphanedVersionCandidate {
|
||||||
qdrantCollection: string;
|
qdrantCollection: string;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
export interface PendingOcrIdentity {
|
||||||
|
sourceId: string;
|
||||||
|
originalManifestHash: string;
|
||||||
|
processingFingerprint: string;
|
||||||
|
metadataHash: string;
|
||||||
|
}
|
||||||
|
|
||||||
|
export interface OcrJobRow {
|
||||||
|
jobId: string;
|
||||||
|
documentId: string;
|
||||||
|
remoteJobId: string | null;
|
||||||
|
remoteIdempotencyKey: string;
|
||||||
|
state: "queued" | "running" | "succeeded" | "failed";
|
||||||
|
requestedPages: number[];
|
||||||
|
completedPages: number;
|
||||||
|
configVersion: string;
|
||||||
|
attemptCount: number;
|
||||||
|
heartbeatAt: Date | null;
|
||||||
|
leaseExpiresAt: Date | null;
|
||||||
|
nextAttemptAt: Date | null;
|
||||||
|
errorCode: string | null;
|
||||||
|
errorDetail: string | null;
|
||||||
|
}
|
||||||
|
|
||||||
|
export interface PendingOcrVersion {
|
||||||
|
version: CatalogVersionRow;
|
||||||
|
jobs: OcrJobRow[];
|
||||||
|
}
|
||||||
|
|
||||||
type VersionDbRow = {
|
type VersionDbRow = {
|
||||||
version_id: string;
|
version_id: string;
|
||||||
source_id: string;
|
source_id: string;
|
||||||
|
|
@ -89,6 +118,37 @@ type VersionDbRow = {
|
||||||
qdrant_collection: string;
|
qdrant_collection: string;
|
||||||
};
|
};
|
||||||
|
|
||||||
|
type OcrJobDbRow = {
|
||||||
|
ocr_job_id: string;
|
||||||
|
ocr_document_id: string;
|
||||||
|
ocr_remote_job_id: string | null;
|
||||||
|
ocr_remote_idempotency_key: string;
|
||||||
|
ocr_state: OcrJobRow["state"];
|
||||||
|
ocr_requested_pages: number[];
|
||||||
|
ocr_completed_pages: string | number;
|
||||||
|
ocr_config_version: string;
|
||||||
|
ocr_attempt_count: string | number;
|
||||||
|
ocr_heartbeat_at: Date | null;
|
||||||
|
ocr_lease_expires_at: Date | null;
|
||||||
|
ocr_next_attempt_at: Date | null;
|
||||||
|
ocr_error_code: string | null;
|
||||||
|
ocr_error_detail: string | null;
|
||||||
|
};
|
||||||
|
|
||||||
|
const OCR_JOB_COLUMNS = [
|
||||||
|
["job_id", "ocr_job_id"], ["document_id", "ocr_document_id"],
|
||||||
|
["remote_job_id", "ocr_remote_job_id"], ["remote_idempotency_key", "ocr_remote_idempotency_key"],
|
||||||
|
["state", "ocr_state"], ["requested_pages", "ocr_requested_pages"],
|
||||||
|
["completed_pages", "ocr_completed_pages"], ["config_version", "ocr_config_version"],
|
||||||
|
["attempt_count", "ocr_attempt_count"], ["heartbeat_at", "ocr_heartbeat_at"],
|
||||||
|
["lease_expires_at", "ocr_lease_expires_at"], ["next_attempt_at", "ocr_next_attempt_at"],
|
||||||
|
["error_code", "ocr_error_code"], ["error_detail", "ocr_error_detail"]
|
||||||
|
] as const;
|
||||||
|
|
||||||
|
function ocrJobColumns(prefix = ""): string {
|
||||||
|
return OCR_JOB_COLUMNS.map(([column, alias]) => `${prefix}${column} AS ${alias}`).join(", ");
|
||||||
|
}
|
||||||
|
|
||||||
function toVersion(row: VersionDbRow): CatalogVersionRow {
|
function toVersion(row: VersionDbRow): CatalogVersionRow {
|
||||||
return {
|
return {
|
||||||
versionId: row.version_id,
|
versionId: row.version_id,
|
||||||
|
|
@ -110,6 +170,25 @@ function toVersion(row: VersionDbRow): CatalogVersionRow {
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
|
function toOcrJob(row: OcrJobDbRow): OcrJobRow {
|
||||||
|
return {
|
||||||
|
jobId: row.ocr_job_id,
|
||||||
|
documentId: row.ocr_document_id,
|
||||||
|
remoteJobId: row.ocr_remote_job_id,
|
||||||
|
remoteIdempotencyKey: row.ocr_remote_idempotency_key,
|
||||||
|
state: row.ocr_state,
|
||||||
|
requestedPages: row.ocr_requested_pages,
|
||||||
|
completedPages: Number(row.ocr_completed_pages),
|
||||||
|
configVersion: row.ocr_config_version,
|
||||||
|
attemptCount: Number(row.ocr_attempt_count),
|
||||||
|
heartbeatAt: row.ocr_heartbeat_at,
|
||||||
|
leaseExpiresAt: row.ocr_lease_expires_at,
|
||||||
|
nextAttemptAt: row.ocr_next_attempt_at,
|
||||||
|
errorCode: row.ocr_error_code,
|
||||||
|
errorDetail: row.ocr_error_detail
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
export class CatalogRepository {
|
export class CatalogRepository {
|
||||||
constructor(private readonly pool: PgPool) {}
|
constructor(private readonly pool: PgPool) {}
|
||||||
|
|
||||||
|
|
@ -220,6 +299,93 @@ export class CatalogRepository {
|
||||||
return result.rows[0] ? toVersion(result.rows[0]) : undefined;
|
return result.rows[0] ? toVersion(result.rows[0]) : undefined;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async findPendingOcrVersion(input: PendingOcrIdentity): Promise<PendingOcrVersion | undefined> {
|
||||||
|
const result = await this.pool.query<VersionDbRow & OcrJobDbRow>(
|
||||||
|
`SELECT v.*, ${ocrJobColumns("j.")}
|
||||||
|
FROM rag_source_versions v
|
||||||
|
JOIN rag_version_documents d ON d.version_id = v.version_id
|
||||||
|
JOIN rag_ocr_jobs j ON j.version_id = d.version_id AND j.document_id = d.document_id
|
||||||
|
WHERE v.source_id = $1
|
||||||
|
AND v.original_manifest_hash = $2
|
||||||
|
AND v.processing_fingerprint = $3
|
||||||
|
AND v.metadata_hash = $4
|
||||||
|
AND v.source_content_hash IS NULL
|
||||||
|
AND v.state IN ('pending', 'indexing', 'review_required')
|
||||||
|
ORDER BY j.document_id`,
|
||||||
|
[input.sourceId, input.originalManifestHash, input.processingFingerprint, input.metadataHash]
|
||||||
|
);
|
||||||
|
return result.rows[0] ? { version: toVersion(result.rows[0]), jobs: result.rows.map(toOcrJob) } : undefined;
|
||||||
|
}
|
||||||
|
|
||||||
|
async claimNextOcrJob(leaseMs: number): Promise<OcrJobRow | undefined> {
|
||||||
|
return withTransaction(this.pool, async (client) => {
|
||||||
|
const candidate = await client.query<{ job_id: string }>(
|
||||||
|
`SELECT job_id FROM rag_ocr_jobs
|
||||||
|
WHERE state = 'queued' AND (next_attempt_at IS NULL OR next_attempt_at <= now())
|
||||||
|
ORDER BY created_at FOR UPDATE SKIP LOCKED LIMIT 1`
|
||||||
|
);
|
||||||
|
if (!candidate.rowCount) return undefined;
|
||||||
|
const result = await client.query<OcrJobDbRow>(
|
||||||
|
`UPDATE rag_ocr_jobs
|
||||||
|
SET state = 'running', attempt_count = attempt_count + 1,
|
||||||
|
started_at = COALESCE(started_at, now()), heartbeat_at = now(),
|
||||||
|
lease_expires_at = now() + ($2::text || ' milliseconds')::interval
|
||||||
|
WHERE job_id = $1
|
||||||
|
RETURNING ${ocrJobColumns()}`,
|
||||||
|
[candidate.rows[0].job_id, String(leaseMs)]
|
||||||
|
);
|
||||||
|
return result.rows[0] ? toOcrJob(result.rows[0]) : undefined;
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
async recoverExpiredOcrLeases(): Promise<OcrJobRow[]> {
|
||||||
|
const result = await this.pool.query<OcrJobDbRow>(
|
||||||
|
`UPDATE rag_ocr_jobs
|
||||||
|
SET state = 'queued', heartbeat_at = NULL, lease_expires_at = NULL, next_attempt_at = now()
|
||||||
|
WHERE state = 'running' AND lease_expires_at <= now()
|
||||||
|
RETURNING ${ocrJobColumns()}`
|
||||||
|
);
|
||||||
|
return result.rows.map(toOcrJob);
|
||||||
|
}
|
||||||
|
|
||||||
|
async markReviewRequired(versionId: string): Promise<void> {
|
||||||
|
const result = await this.pool.query(
|
||||||
|
`UPDATE rag_source_versions SET state = 'review_required'
|
||||||
|
WHERE version_id = $1 AND state = 'indexing'
|
||||||
|
AND EXISTS (SELECT 1 FROM rag_ocr_jobs WHERE version_id = $1)
|
||||||
|
AND NOT EXISTS (SELECT 1 FROM rag_ocr_jobs WHERE version_id = $1 AND state <> 'succeeded')
|
||||||
|
AND NOT EXISTS (
|
||||||
|
SELECT 1 FROM rag_version_documents d
|
||||||
|
WHERE d.version_id = $1 AND d.content_hash IS NULL
|
||||||
|
AND NOT EXISTS (
|
||||||
|
SELECT 1 FROM rag_ocr_jobs j
|
||||||
|
WHERE j.version_id = d.version_id AND j.document_id = d.document_id AND j.state = 'succeeded'
|
||||||
|
)
|
||||||
|
)`,
|
||||||
|
[versionId]
|
||||||
|
);
|
||||||
|
if (result.rowCount !== 1) throw new CatalogError("Version is not ready for OCR review", 409, "INVALID_VERSION_STATE");
|
||||||
|
}
|
||||||
|
|
||||||
|
async markOcrReviewIndexing(versionId: string, reviewedBy: string): Promise<void> {
|
||||||
|
const result = await this.pool.query(
|
||||||
|
`UPDATE rag_source_versions SET state = 'indexing', reviewed_at = now(), reviewed_by = $2, indexing_started_at = now()
|
||||||
|
WHERE version_id = $1 AND state = 'review_required'`,
|
||||||
|
[versionId, reviewedBy]
|
||||||
|
);
|
||||||
|
if (result.rowCount !== 1) throw new CatalogError("Version is not awaiting OCR review", 409, "INVALID_VERSION_STATE");
|
||||||
|
}
|
||||||
|
|
||||||
|
async rejectOcrVersion(versionId: string, reviewedBy: string, reason: string): Promise<void> {
|
||||||
|
const result = await this.pool.query(
|
||||||
|
`UPDATE rag_source_versions
|
||||||
|
SET state = 'rejected', reviewed_at = now(), reviewed_by = $2, error_code = 'OCR_REJECTED', error_detail = $3
|
||||||
|
WHERE version_id = $1 AND state = 'review_required'`,
|
||||||
|
[versionId, reviewedBy, reason.slice(0, 2000)]
|
||||||
|
);
|
||||||
|
if (result.rowCount !== 1) throw new CatalogError("Version is not awaiting OCR review", 409, "INVALID_VERSION_STATE");
|
||||||
|
}
|
||||||
|
|
||||||
async createPendingVersion(input: CreateVersionInput): Promise<CatalogVersionRow> {
|
async createPendingVersion(input: CreateVersionInput): Promise<CatalogVersionRow> {
|
||||||
return withTransaction(this.pool, async (client) => {
|
return withTransaction(this.pool, async (client) => {
|
||||||
if (await this.isMaintenanceEnabled(client, true)) {
|
if (await this.isMaintenanceEnabled(client, true)) {
|
||||||
|
|
|
||||||
147
tests/catalog/repository-ocr.test.ts
Normal file
147
tests/catalog/repository-ocr.test.ts
Normal file
|
|
@ -0,0 +1,147 @@
|
||||||
|
import assert from "node:assert/strict";
|
||||||
|
import test from "node:test";
|
||||||
|
import { CatalogError, CatalogRepository } from "../../src/modules/catalog/repository.js";
|
||||||
|
|
||||||
|
const versionRow = {
|
||||||
|
version_id: "version-1",
|
||||||
|
source_id: "source-1",
|
||||||
|
version_number: 2,
|
||||||
|
previous_version_id: "version-0",
|
||||||
|
state: "indexing",
|
||||||
|
tags: [],
|
||||||
|
source_content_hash: null,
|
||||||
|
processing_fingerprint: "fingerprint",
|
||||||
|
metadata_hash: "metadata",
|
||||||
|
embedding_provider: "test",
|
||||||
|
embedding_model: "test",
|
||||||
|
embedding_dimensions: 3,
|
||||||
|
expected_document_count: 1,
|
||||||
|
expected_point_count: 0,
|
||||||
|
verified_point_count: 0,
|
||||||
|
qdrant_collection: "rag"
|
||||||
|
};
|
||||||
|
|
||||||
|
const jobRow = {
|
||||||
|
ocr_job_id: "job-1",
|
||||||
|
ocr_document_id: "document-1",
|
||||||
|
ocr_remote_job_id: "remote-1",
|
||||||
|
ocr_remote_idempotency_key: "document-hash:ocr-v1:pages-hash",
|
||||||
|
ocr_state: "running",
|
||||||
|
ocr_requested_pages: [1, 3],
|
||||||
|
ocr_completed_pages: 1,
|
||||||
|
ocr_config_version: "ocr-v1",
|
||||||
|
ocr_attempt_count: 2,
|
||||||
|
ocr_heartbeat_at: new Date("2026-09-14T10:00:00Z"),
|
||||||
|
ocr_lease_expires_at: new Date("2026-09-14T10:01:00Z"),
|
||||||
|
ocr_next_attempt_at: null,
|
||||||
|
ocr_error_code: null,
|
||||||
|
ocr_error_detail: null
|
||||||
|
};
|
||||||
|
|
||||||
|
function poolFor(handler: (sql: string, params?: unknown[]) => Promise<{ rowCount: number; rows: unknown[] }>) {
|
||||||
|
const client = {
|
||||||
|
query: async (sql: string, params?: unknown[]) => /^(BEGIN|COMMIT|ROLLBACK|SET CONSTRAINTS)/.test(sql.trim())
|
||||||
|
? { rowCount: 0, rows: [] }
|
||||||
|
: handler(sql, params),
|
||||||
|
release() {}
|
||||||
|
};
|
||||||
|
return { query: handler, connect: async () => client };
|
||||||
|
}
|
||||||
|
|
||||||
|
test("pending OCR lookup returns the persisted version and jobs for a duplicate ingestion", async () => {
|
||||||
|
const pool = poolFor(async (sql, params) => {
|
||||||
|
assert.match(sql, /JOIN rag_version_documents/);
|
||||||
|
assert.match(sql, /JOIN rag_ocr_jobs/);
|
||||||
|
assert.match(sql, /source_content_hash IS NULL/);
|
||||||
|
assert.match(sql, /state IN \('pending', 'indexing', 'review_required'\)/);
|
||||||
|
assert.deepEqual(params, ["source-1", "manifest", "fingerprint", "metadata"]);
|
||||||
|
return { rowCount: 1, rows: [{ ...versionRow, ...jobRow }] };
|
||||||
|
});
|
||||||
|
const repository = new CatalogRepository(pool as never);
|
||||||
|
|
||||||
|
const candidate = await repository.findPendingOcrVersion({
|
||||||
|
sourceId: "source-1",
|
||||||
|
originalManifestHash: "manifest",
|
||||||
|
processingFingerprint: "fingerprint",
|
||||||
|
metadataHash: "metadata"
|
||||||
|
});
|
||||||
|
|
||||||
|
assert.equal(candidate?.version.versionId, "version-1");
|
||||||
|
assert.deepEqual(candidate?.jobs.map((job) => [job.jobId, job.remoteJobId, job.remoteIdempotencyKey]), [
|
||||||
|
["job-1", "remote-1", "document-hash:ocr-v1:pages-hash"]
|
||||||
|
]);
|
||||||
|
});
|
||||||
|
|
||||||
|
test("pending OCR lookup returns undefined when only terminal candidates exist", async () => {
|
||||||
|
const repository = new CatalogRepository(poolFor(async () => ({ rowCount: 0, rows: [] })) as never);
|
||||||
|
const candidate = await repository.findPendingOcrVersion({
|
||||||
|
sourceId: "source-1",
|
||||||
|
originalManifestHash: "manifest",
|
||||||
|
processingFingerprint: "fingerprint",
|
||||||
|
metadataHash: "metadata"
|
||||||
|
});
|
||||||
|
assert.equal(candidate, undefined);
|
||||||
|
});
|
||||||
|
|
||||||
|
test("lease claiming uses a locked queue row and keeps its remote identity", async () => {
|
||||||
|
const calls: string[] = [];
|
||||||
|
const pool = poolFor(async (sql, params) => {
|
||||||
|
calls.push(sql);
|
||||||
|
if (sql.includes("FOR UPDATE SKIP LOCKED")) {
|
||||||
|
assert.match(sql, /state = 'queued'/);
|
||||||
|
return { rowCount: 1, rows: [{ job_id: "job-1" }] };
|
||||||
|
}
|
||||||
|
assert.match(sql, /attempt_count = attempt_count \+ 1/);
|
||||||
|
assert.deepEqual(params, ["job-1", "30000"]);
|
||||||
|
return { rowCount: 1, rows: [jobRow] };
|
||||||
|
});
|
||||||
|
const repository = new CatalogRepository(pool as never);
|
||||||
|
|
||||||
|
const claimed = await repository.claimNextOcrJob(30_000);
|
||||||
|
|
||||||
|
assert.equal(calls.length, 2);
|
||||||
|
assert.equal(claimed?.remoteJobId, "remote-1");
|
||||||
|
assert.equal(claimed?.remoteIdempotencyKey, "document-hash:ocr-v1:pages-hash");
|
||||||
|
});
|
||||||
|
|
||||||
|
test("lease claiming returns undefined when no queued job is eligible", async () => {
|
||||||
|
const repository = new CatalogRepository(poolFor(async () => ({ rowCount: 0, rows: [] })) as never);
|
||||||
|
assert.equal(await repository.claimNextOcrJob(30_000), undefined);
|
||||||
|
});
|
||||||
|
|
||||||
|
test("expired leases return to queued without clearing remote recovery identity", async () => {
|
||||||
|
const repository = new CatalogRepository(poolFor(async (sql) => {
|
||||||
|
assert.match(sql, /state = 'running'/);
|
||||||
|
assert.match(sql, /lease_expires_at <= now\(\)/);
|
||||||
|
assert.doesNotMatch(sql, /remote_(job_id|idempotency_key)\s*=/);
|
||||||
|
return { rowCount: 1, rows: [{ ...jobRow, ocr_state: "queued" }] };
|
||||||
|
}) as never);
|
||||||
|
|
||||||
|
const recovered = await repository.recoverExpiredOcrLeases();
|
||||||
|
|
||||||
|
assert.deepEqual(recovered.map((job) => [job.state, job.remoteJobId, job.remoteIdempotencyKey]), [
|
||||||
|
["queued", "remote-1", "document-hash:ocr-v1:pages-hash"]
|
||||||
|
]);
|
||||||
|
});
|
||||||
|
|
||||||
|
test("review transitions enforce indexing, approval, and rejection state guards", async () => {
|
||||||
|
const successfulSql: string[] = [];
|
||||||
|
const success = new CatalogRepository(poolFor(async (sql) => {
|
||||||
|
successfulSql.push(sql);
|
||||||
|
return { rowCount: 1, rows: [] };
|
||||||
|
}) as never);
|
||||||
|
await success.markReviewRequired("version-1");
|
||||||
|
await success.markOcrReviewIndexing("version-1", "reviewer-1");
|
||||||
|
await success.rejectOcrVersion("version-1", "reviewer-1", "Unreadable code");
|
||||||
|
assert.match(successfulSql[0], /state = 'indexing'.*NOT EXISTS/s);
|
||||||
|
assert.match(successfulSql[0], /FROM rag_version_documents/);
|
||||||
|
assert.match(successfulSql[0], /rag_ocr_jobs WHERE version_id = \$1 AND state <> 'succeeded'/);
|
||||||
|
assert.match(successfulSql[1], /state = 'review_required'/);
|
||||||
|
assert.match(successfulSql[2], /state = 'review_required'/);
|
||||||
|
|
||||||
|
const blocked = new CatalogRepository(poolFor(async () => ({ rowCount: 0, rows: [] })) as never);
|
||||||
|
await assert.rejects(
|
||||||
|
blocked.markReviewRequired("version-2"),
|
||||||
|
(error) => error instanceof CatalogError && error.code === "INVALID_VERSION_STATE"
|
||||||
|
);
|
||||||
|
});
|
||||||
Loading…
Add table
Reference in a new issue