/** * Entity match workbench: pending candidates with BOTH facility rows side by side (name, operator, city, address, * coordinates + distance, facility codes, external ids, source counts) and the matcher's score / reasons. * Decisions: approve (merge), reject (keep separate), related-campus (candidate is a building of the matched campus), * defer. */ import type { FastifyInstance } from "fastify"; import { z } from "zod"; import { extractFacilityCodes } from "@dci/core"; import { envelope, notFound, parseQuery, parseBody, badRequest } from "../../lib/http.js"; import { intParam, pageParam, strParam } from "../../lib/params.js"; import { pg, page } from "../../lib/sql.js"; import { int, iso, json, num, record, reqStr, str, strArray, type Row } from "../../lib/rows.js"; import { facilitiesByIds } from "../../repositories/facilities.js"; import { ensureManualSource, mergeFacility, setFacilityParent } from "../../repositories/admin/merge.js"; import { invalidate } from "../../cache.js"; const STATUSES = ["pending", "approved", "rejected", "related_campus", "deferred", "auto_merged", "auto_created", "all"] as const; function matchDto(r: Row): Record { return { id: reqStr(r.id), connectorId: reqStr(r.connector_id), candidateKey: reqStr(r.candidate_key), candidate: json(r.candidate, {}), matchedFacilityId: str(r.matched_facility_id), score: num(r.score), reasons: strArray(r.reasons), status: reqStr(r.status), decidedBy: str(r.decided_by), decidedAt: iso(r.decided_at), createdAt: iso(r.created_at), candidateFacilityId: str(r.candidate_facility_id) }; } interface Side { id: string; slug: string; name: string; operator: string | null; city: string | null; address: string | null; lat: number | null; lng: number | null; codes: string[]; externalIds: Record; sourceCount: number; status: string; recordScope: string; parentFacilityId: string | null } function side(r: Row): Side { return { id: reqStr(r.id), slug: reqStr(r.slug), name: reqStr(r.name), operator: str(r.op_name), city: str(r.city), address: str(r.address), lat: num(r.lat), lng: num(r.lng), codes: extractFacilityCodes(reqStr(r.name)), externalIds: record(r.external_ids), sourceCount: int(r.source_count), status: reqStr(r.status), recordScope: reqStr(r.record_scope, "facility"), parentFacilityId: str(r.parent_facility_id) }; } function haversineKm(a: Side, b: Side): number | null { if (a.lat == null || a.lng == null || b.lat == null || b.lng == null) return null; const R = 6371.0088, toRad = (d: number) => (d * Math.PI) / 180; const dLat = toRad(b.lat - a.lat), dLng = toRad(b.lng - a.lng); const h = Math.sin(dLat / 2) ** 2 + Math.cos(toRad(a.lat)) * Math.cos(toRad(b.lat)) * Math.sin(dLng / 2) ** 2; return Math.round(2 * R * Math.asin(Math.sqrt(Math.min(1, h))) * 1000) / 1000; } /** The facility created from the candidate (entity key or candidate.createdFacilityId) — for a pending row. */ function candidateFacilityId(r: Row): string | null { return str(r.candidate_facility_id) ?? str((json>(r.candidate, {}) as Record).createdFacilityId); } export async function matchAdminRoutes(app: FastifyInstance): Promise { app.get("/matches", { schema: { summary: "Entity match queue (?status=pending|approved|rejected|related_campus|deferred|auto_merged|auto_created|all&connector=) with both facility rows (pair: name, operator, city, address, coords + distance, codes, external ids, source counts), score and reasons" } }, async (req) => { const q = parseQuery(z.object({ status: z.preprocess((v) => (v === "" || v == null ? undefined : v), z.enum(STATUSES).optional()), connector: strParam, page: pageParam, per_page: intParam }), req.query); const sql = pg(); const pg_ = page(q.page, q.per_page, 200, 50); const status = q.status ?? "pending"; const rows = await sql` select m.*, coalesce(k.entity_id, m.candidate->>'createdFacilityId') as candidate_facility_id, count(*) over() as total from entity_matches m left join entity_keys k on k.key = m.candidate_key and k.entity_type = 'facility' where ${status === "all" ? sql`true` : sql`m.status = ${status}`} and ${q.connector ? sql`m.connector_id = ${q.connector}` : sql`true`} order by m.score desc, m.created_at desc limit ${pg_.perPage} offset ${pg_.offset}`; const ids = [...new Set(rows.flatMap((r) => [str(r.matched_facility_id), candidateFacilityId(r)]).filter((x): x is string => Boolean(x)))]; const [facs, sides] = await Promise.all([ facilitiesByIds(ids), ids.length ? sql`select f.id, f.slug, f.name, f.city, f.address, f.lat, f.lng, f.external_ids, f.source_count, f.status, f.record_scope, f.parent_facility_id, o.name as op_name from facilities f left join operators o on o.id = f.operator_id where f.id = any(${ids})` : Promise.resolve([] as Row[]), ]); const by = new Map(facs.map((f) => [f.id, f])); const sideBy = new Map(sides.map((r) => [reqStr(r.id), side(r)])); const items = rows.map((r) => { const matchedId = str(r.matched_facility_id); const candId = candidateFacilityId(r); const cand = json>(r.candidate, {}); const a = candId ? sideBy.get(candId) ?? null : null; const b = matchedId ? sideBy.get(matchedId) ?? null : null; // when the candidate was never persisted as a facility, describe it from the candidate payload itself const candidateSide: Partial | null = a ?? (Object.keys(cand).length ? { id: "", slug: "", name: str(cand.name) ?? "", operator: str(cand.operatorName ?? cand.operator), city: str(cand.city), address: str(cand.address), lat: num((cand.geo as Record | undefined)?.lat ?? cand.lat), lng: num((cand.geo as Record | undefined)?.lng ?? cand.lng), codes: extractFacilityCodes(str(cand.name) ?? ""), externalIds: record(cand.externalIds), sourceCount: 0, status: str(cand.status) ?? "unknown", recordScope: "facility", parentFacilityId: null } : null); const distanceKm = candidateSide && b && candidateSide.lat != null && candidateSide.lng != null ? haversineKm(candidateSide as Side, b) : null; return { ...matchDto(r), candidateFacilityId: candId, matched: matchedId ? by.get(matchedId) ?? null : null, candidateFacility: candId ? by.get(candId) ?? null : null, pair: { candidate: candidateSide, matched: b, distanceKm, sameCodes: candidateSide && b ? candidateSide.codes!.filter((c) => b.codes.includes(c)) : [] }, }; }); return envelope(items, { total: rows.length ? int(rows[0]!.total) : 0, page: pg_.page, perPage: pg_.perPage }); }); async function pendingMatch(id: string): Promise { const sql = pg(); const rows = await sql`select m.*, coalesce(k.entity_id, m.candidate->>'createdFacilityId') as candidate_facility_id from entity_matches m left join entity_keys k on k.key = m.candidate_key and k.entity_type = 'facility' where m.id = ${id}`; const m = rows[0]; if (!m) throw notFound("match"); if (m.status !== "pending" && m.status !== "deferred") throw badRequest(`match is already ${String(m.status)}`); return m; } app.post("/matches/:id/approve", { schema: { summary: "Approve: point the candidate key at the matched facility and merge any separately-created facility into it" } }, async (req) => { const { id } = req.params as { id: string }; const sql = pg(); const m = await pendingMatch(id); const matched = str(m.matched_facility_id); if (!matched) throw badRequest("match has no matched_facility_id"); const target = await sql`select id, merged_into from facilities where id = ${matched}`; if (!target[0]) throw badRequest("matched facility no longer exists"); const survivor = str(target[0].merged_into) ?? matched; await ensureManualSource(); const key = reqStr(m.candidate_key); const existing = await sql`select entity_id from entity_keys where key = ${key} and entity_type = 'facility'`; const separate = str(existing[0]?.entity_id) ?? candidateFacilityId(m); let merge: Awaited> | null = null; if (separate && separate !== survivor) merge = await mergeFacility(separate, survivor); await sql`insert into entity_keys (key, entity_type, entity_id, connector_id) values (${key}, 'facility', ${survivor}, ${reqStr(m.connector_id)}) on conflict (key, entity_type) do update set entity_id = ${survivor}`; await sql`update entity_matches set status = 'approved', decided_by = 'admin', decided_at = now() where id = ${id}`; await invalidate("/datacenters"); return envelope({ id, status: "approved", facilityId: survivor, candidateKey: key, merged: merge }); }); app.post("/matches/:id/reject", { schema: { summary: "Reject: keep the candidate as a separate facility" } }, async (req) => { const { id } = req.params as { id: string }; const sql = pg(); const rows = await sql`update entity_matches set status = 'rejected', decided_by = 'admin', decided_at = now() where id = ${id} and status in ('pending', 'deferred') returning id`; if (!rows[0]) { const exists = await sql`select status from entity_matches where id = ${id}`; if (!exists.length) throw notFound("match"); throw badRequest(`match is already ${String(exists[0]!.status)}`); } return envelope({ id, status: "rejected" }); }); app.post("/matches/:id/related-campus", { schema: { summary: "Related campus: the candidate facility is a BUILDING of the matched campus (parent_facility_id set, record_scope building / campus); both rows stay, aggregates never count both" } }, async (req) => { const { id } = req.params as { id: string }; const sql = pg(); const m = await pendingMatch(id); const matched = str(m.matched_facility_id); const building = candidateFacilityId(m); if (!matched) throw badRequest("match has no matched_facility_id"); if (!building) throw badRequest("the candidate was never persisted as a facility — approve or reject instead"); const parentRow = (await sql`select id, merged_into from facilities where id = ${matched}`)[0]; if (!parentRow) throw badRequest("matched facility no longer exists"); const campus = str(parentRow.merged_into) ?? matched; const res = await setFacilityParent(building, campus); await sql`update entity_matches set status = 'related_campus', decided_by = 'admin', decided_at = now() where id = ${id}`; await Promise.all([invalidate("/datacenters"), invalidate("/dashboard"), invalidate("/map")]); return envelope({ id, status: "related_campus", building, campus, containment: res }); }); app.post("/matches/:id/defer", { schema: { summary: "Defer the decision (status deferred, keeps the row in the workbench under ?status=deferred)" } }, async (req) => { const { id } = req.params as { id: string }; const sql = pg(); const rows = await sql`update entity_matches set status = 'deferred', decided_by = 'admin', decided_at = now() where id = ${id} and status = 'pending' returning id`; if (!rows[0]) { const exists = await sql`select status from entity_matches where id = ${id}`; if (!exists.length) throw notFound("match"); throw badRequest(`match is already ${String(exists[0]!.status)}`); } return envelope({ id, status: "deferred" }); }); }