/** entity_keys helpers (connector-scoped stable keys → entity ids), with merged_into following for facilities. */ import { sql } from "@dci/db"; import type { IngestContext, Tx } from "./common.js"; export async function entityIdForKey(tx: Tx, key: string, entityType: string): Promise { const r = await tx.execute(sql`select entity_id from entity_keys where key = ${key} and entity_type = ${entityType} limit 1`); return r[0] ? String(r[0].entity_id) : null; } /** Facility id for a connector key, following `merged_into` chains (max 5 hops). Cached per batch. */ export async function facilityIdForKey(tx: Tx, ctx: IngestContext, key: string): Promise { const cached = ctx.caches.facilityIdByKey.get(key); if (cached) return cached; let id: string | null = await entityIdForKey(tx, key, "facility"); if (!id) return null; for (let hop = 0; hop < 5; hop++) { const rows: Array> = await tx.execute(sql`select merged_into from facilities where id = ${id}`); const m: unknown = rows[0]?.merged_into; if (!m) break; id = String(m); } ctx.caches.facilityIdByKey.set(key, id); return id; } export async function upsertKey(tx: Tx, ctx: IngestContext, key: string, entityType: string, entityId: string): Promise { await tx.execute(sql`insert into entity_keys (key, entity_type, entity_id, connector_id) values (${key}, ${entityType}, ${entityId}, ${ctx.run.connectorId}) on conflict (key, entity_type) do update set entity_id = excluded.entity_id, connector_id = excluded.connector_id`); if (entityType === "facility") ctx.caches.facilityIdByKey.set(key, entityId); }