SPB Git forge
38commits 1branches 0releases
338.7 MBsize
maindefault branch
5 h agolast push
HTML 53.9% TypeScript 44.5% JavaScript 0.6% SQL 0.5%
9.0 KB · 121 lines typescript
Raw Blame History
1import type { FastifyInstance } from "fastify";2import { z } from "zod";3import { envelope, notFound, parseBody, parseQuery } from "../../lib/http.js";4import { intParam, pageParam, strParam } from "../../lib/params.js";5import { pg, page, andAll, likePattern, type Fragment } from "../../lib/sql.js";6import { int, iso, json, num, reqStr, str, type Row } from "../../lib/rows.js";7import { resolveEntityRefs } from "../../lib/resolve.js";8import { enqueueCrawl, manualJobId } from "../../queues.js";9import { getRawObject } from "../../storage.js";10import { provenanceDto } from "../../lib/dto.js";1112function docDto(r: Row): Record<string, unknown> {13  return {14    id: reqStr(r.id), connectorId: reqStr(r.connector_id), sourceId: reqStr(r.source_id), url: reqStr(r.url), canonicalUrl: reqStr(r.canonical_url), pageType: reqStr(r.page_type), classifier: str(r.classifier), fetchLevel: int(r.fetch_level, 1), priority: int(r.priority, 50),15    contentHash: str(r.content_hash), etag: str(r.etag), lastModified: str(r.last_modified), statusCode: num(r.status_code), contentType: str(r.content_type), sizeBytes: num(r.size_bytes), title: str(r.title), storageKey: str(r.storage_key), extractorVersion: str(r.extractor_version),16    extractOk: r.extract_ok == null ? null : r.extract_ok === true, extractCount: int(r.extract_count), error: str(r.error), errorCount: int(r.error_count), firstSeen: iso(r.first_seen), lastFetched: iso(r.last_fetched), lastChanged: iso(r.last_changed), lastChecked: iso(r.last_checked), nextCheck: iso(r.next_check),17    changeFrequencyScore: num(r.change_frequency_score), fetchCount: int(r.fetch_count), changeCount: int(r.change_count), discoveredFrom: str(r.discovered_from), quarantined: r.quarantined === true, entityRefs: json<Array<{ type: string; id: string }>>(r.entity_refs, []), updatedAt: iso(r.updated_at),18    versionCount: r.version_count != null ? int(r.version_count) : undefined,19  };20}2122function versionDto(v: Row): Record<string, unknown> {23  return { id: reqStr(v.id), documentId: reqStr(v.document_id), contentHash: reqStr(v.content_hash), fetchedAt: iso(v.fetched_at), fetchLevel: int(v.fetch_level, 1), statusCode: num(v.status_code), sizeBytes: num(v.size_bytes), storageKey: str(v.storage_key), significance: int(v.significance), diffSummary: json(v.diff_summary, null), detectedChanges: json(v.detected_changes, []) };24}2526export async function documentAdminRoutes(app: FastifyInstance): Promise<void> {27  app.get("/documents", { schema: { summary: "Documents (?connector=&pageType=&status=error|changed|quarantined&q=&page=)" } }, async (req) => {28    const q = parseQuery(z.object({ connector: strParam, pageType: strParam, status: strParam, q: strParam, page: pageParam, per_page: intParam }), req.query);29    const sql = pg();30    const pg_ = page(q.page, q.per_page, 200, 50);31    const c: Fragment[] = [];32    if (q.connector) c.push(sql`d.connector_id = ${q.connector}`);33    if (q.pageType) c.push(sql`d.page_type = ${q.pageType}`);34    if (q.status === "error") c.push(sql`d.error is not null`);35    if (q.status === "changed") c.push(sql`d.change_count > 0`);36    if (q.status === "quarantined") c.push(sql`d.quarantined`);37    if (q.q) c.push(sql`(d.url ilike ${likePattern(q.q)} or d.title ilike ${likePattern(q.q)})`);38    const rows = await sql<Row[]>`select d.*, count(*) over() as total from documents d where ${andAll(sql, c)} order by coalesce(d.last_changed, d.last_fetched, d.first_seen) desc limit ${pg_.perPage} offset ${pg_.offset}`;39    return envelope(rows.map(docDto), { total: rows.length ? int(rows[0]!.total) : 0, page: pg_.page, perPage: pg_.perPage });40  });4142  app.get("/documents/:id", { schema: { summary: "Document row + versions + provenance referencing it + resolved entity refs" } }, async (req) => {43    const { id } = req.params as { id: string };44    const sql = pg();45    const rows = await sql<Row[]>`select d.*, (select count(*)::int from document_versions v where v.document_id = d.id) as version_count from documents d where d.id = ${id}`;46    const d = rows[0];47    if (!d) throw notFound("document");48    const [versions, prov] = await Promise.all([49      sql<Row[]>`select * from document_versions where document_id = ${id} order by fetched_at desc limit 100`,50      sql<Row[]>`select p.*, s.name as source_name, s.kind as source_kind from provenance p left join sources s on s.id = p.source_id where p.document_id = ${id} order by p.entity_type, p.entity_id, p.field limit 500`,51    ]);52    const refs = json<Array<{ type: string; id: string }>>(d.entity_refs, []);53    const resolved = await resolveEntityRefs(refs);54    return envelope({55      document: docDto(d),56      versions: versions.map(versionDto),57      provenance: prov.map((p) => ({ entityType: reqStr(p.entity_type), entityId: reqStr(p.entity_id), isCurrent: p.is_current === true, ...provenanceDto(p) })),58      entities: refs.map((r) => ({ ...r, ...(resolved.get(`${r.type}:${r.id}`) ?? { slug: null, name: null }) })),59    });60  });6162  app.get("/documents/:id/raw", { schema: { summary: "Decompressed raw body from MinIO (content-type preserved, capped 5 MB)" } }, async (req, reply) => {63    const { id } = req.params as { id: string };64    const q = parseQuery(z.object({ version: strParam }), req.query);65    const sql = pg();66    const rows = await sql<Row[]>`select storage_key, content_type from documents where id = ${id}`;67    const d = rows[0];68    if (!d) throw notFound("document");69    let key = str(d.storage_key);70    if (q.version) { const v = await sql<Row[]>`select storage_key from document_versions where id = ${q.version} and document_id = ${id}`; key = str(v[0]?.storage_key) ?? key; }71    if (!key) throw notFound("raw body (no storage key)");72    const obj = await getRawObject(key);73    if (!obj) throw notFound("raw object in storage");74    reply.header("content-type", str(d.content_type) ?? obj.contentType ?? "application/octet-stream");75    reply.header("x-storage-key", key);76    reply.header("x-truncated", obj.truncated ? "1" : "0");77    reply.header("content-disposition", `inline; filename="${id}.${(str(d.content_type) ?? "").includes("pdf") ? "pdf" : (str(d.content_type) ?? "").includes("json") ? "json" : "html"}"`);78    return reply.send(obj.body);79  });8081  app.get("/documents/:id/versions/:vid/diff", { schema: { summary: "Version diff summary + detected field changes" } }, async (req) => {82    const { id, vid } = req.params as { id: string; vid: string };83    const sql = pg();84    const rows = await sql<Row[]>`select * from document_versions where id = ${vid} and document_id = ${id}`;85    const v = rows[0];86    if (!v) throw notFound("document version");87    const prev = await sql<Row[]>`select id, content_hash, fetched_at from document_versions where document_id = ${id} and fetched_at < ${String(v.fetched_at)} order by fetched_at desc limit 1`;88    return envelope({ version: versionDto(v), previous: prev[0] ? { id: reqStr(prev[0].id), contentHash: reqStr(prev[0].content_hash), fetchedAt: iso(prev[0].fetched_at) } : null, diffSummary: json(v.diff_summary, null), detectedChanges: json(v.detected_changes, []) });89  });9091  async function docForJob(id: string): Promise<{ connectorId: string; url: string }> {92    const sql = pg();93    const rows = await sql<Row[]>`select connector_id, url from documents where id = ${id}`;94    if (!rows[0]) throw notFound("document");95    return { connectorId: reqStr(rows[0].connector_id), url: reqStr(rows[0].url) };96  }9798  app.post("/documents/:id/reprocess", { schema: { summary: "Re-run extraction on the stored body (enqueue task reprocess)" } }, async (req) => {99    const { id } = req.params as { id: string };100    const d = await docForJob(id);101    const job = await enqueueCrawl({ connectorId: d.connectorId, task: "reprocess", urls: [d.url], requestedBy: "admin" }, { jobId: manualJobId(d.connectorId, `reprocess-${id}`) });102    return envelope({ enqueued: true, job, documentId: id, url: d.url });103  });104105  app.post("/documents/:id/refetch", { schema: { summary: "Force a fresh fetch of the URL (enqueue crawl with urls + force)" } }, async (req) => {106    const { id } = req.params as { id: string };107    const d = await docForJob(id);108    const job = await enqueueCrawl({ connectorId: d.connectorId, task: "crawl", urls: [d.url], force: true, requestedBy: "admin" }, { jobId: manualJobId(d.connectorId, `refetch-${id}`) });109    return envelope({ enqueued: true, job, documentId: id, url: d.url });110  });111112  app.post("/documents/:id/quarantine", { schema: { summary: "Set / clear quarantine {quarantined: bool}" } }, async (req) => {113    const { id } = req.params as { id: string };114    const body = parseBody(z.object({ quarantined: z.boolean() }), req.body);115    const sql = pg();116    const rows = await sql<Row[]>`update documents set quarantined = ${body.quarantined}, next_check = case when ${body.quarantined} then null else coalesce(next_check, now()) end, updated_at = now() where id = ${id} returning id, quarantined, next_check`;117    if (!rows[0]) throw notFound("document");118    return envelope({ id, quarantined: rows[0].quarantined === true, nextCheck: iso(rows[0].next_check) });119  });120}121