/** * Private watchlists — no accounts. The owner is an opaque random token stored in the httpOnly `dci_watch` cookie; * rows are only ever read / deleted through that token. */ import { randomBytes } from "node:crypto"; import type { EventDTO, WatchlistItem } from "@dci/core"; import { newId } from "@dci/core"; import { pg, eventCols, eventJoins, page } from "../lib/sql.js"; import { reqIso, reqStr, str, type Row } from "../lib/rows.js"; import { findBySlugOrId, findCountry, resolveEntityRefs } from "../lib/resolve.js"; import { toEventDtos } from "./events.js"; export const WATCH_COOKIE = "dci_watch"; export const WATCH_MAX_ITEMS = 200; export const WATCH_ENTITY_TYPES = ["operator", "metro", "country", "project", "facility"] as const; export type WatchEntityType = (typeof WATCH_ENTITY_TYPES)[number]; export function newWatchToken(): string { return randomBytes(24).toString("base64url"); } export function parseCookies(header: string | undefined): Record { const out: Record = {}; if (!header) return out; for (const part of header.split(";")) { const i = part.indexOf("="); if (i < 0) continue; const k = part.slice(0, i).trim(); const v = part.slice(i + 1).trim(); if (k && /^[A-Za-z0-9_-]{8,128}$/.test(v)) out[k] = v; } return out; } const TABLE: Record = { operator: "operators", metro: "metros", project: "projects", facility: "facilities", country: null }; /** Resolve a slug or id to the canonical entity id (null when unknown). */ export async function resolveWatchEntity(type: WatchEntityType, idOrSlug: string): Promise { if (type === "country") { const c = await findCountry(idOrSlug); return c ? reqStr(c.iso2) : null; } const row = await findBySlugOrId(TABLE[type]!, idOrSlug); if (!row) return null; if (type === "project" && (row.hidden === true || row.merged_into)) return str(row.merged_into) ?? null; return reqStr(row.id); } async function toItems(rows: Row[]): Promise { const refs = await resolveEntityRefs(rows.map((r) => ({ type: reqStr(r.entity_type), id: reqStr(r.entity_id) }))); return rows.map((r) => { const ref = refs.get(`${reqStr(r.entity_type)}:${reqStr(r.entity_id)}`); return { id: reqStr(r.id), entityType: reqStr(r.entity_type) as WatchEntityType, entityId: reqStr(r.entity_id), slug: ref?.slug ?? reqStr(r.entity_id), name: ref?.name ?? reqStr(r.entity_id), createdAt: reqIso(r.created_at) }; }); } export async function listWatchlist(token: string): Promise { const sql = pg(); const rows = await sql`select * from watchlists where owner_token = ${token} order by created_at desc limit ${WATCH_MAX_ITEMS}`; return toItems(rows); } export async function addWatch(token: string, type: WatchEntityType, entityId: string): Promise<{ item: WatchlistItem; created: boolean } | { error: "full" }> { const sql = pg(); const n = await sql`select count(*)::int as n from watchlists where owner_token = ${token}`; const existing = await sql`select * from watchlists where owner_token = ${token} and entity_type = ${type} and entity_id = ${entityId}`; if (existing[0]) return { item: (await toItems(existing))[0]!, created: false }; if (Number(n[0]?.n ?? 0) >= WATCH_MAX_ITEMS) return { error: "full" }; const rows = await sql`insert into watchlists (id, owner_token, entity_type, entity_id) values (${newId("watch")}, ${token}, ${type}, ${entityId}) on conflict (owner_token, entity_type, entity_id) do update set entity_id = excluded.entity_id returning *`; return { item: (await toItems(rows))[0]!, created: true }; } export async function removeWatch(token: string, id: string): Promise { const sql = pg(); const rows = await sql`delete from watchlists where owner_token = ${token} and id = ${id} returning id`; return rows.length > 0; } /** Events for the watched entities, newest first, one row per cluster. */ export async function watchFeed(token: string, pageNo?: number, perPage?: number): Promise<{ items: EventDTO[]; total: number; page: number; perPage: number }> { const sql = pg(); const pg_ = page(pageNo, perPage, 100, 50); const rows = await sql` with w as (select entity_type, entity_id from watchlists where owner_token = ${token}), matched as ( select distinct on (coalesce(e.cluster_id, e.id)) e.id from events e where e.review_status <> 'rejected' and ( exists (select 1 from w where w.entity_type = 'operator' and (e.operator_id = w.entity_id or (e.entity_type = 'operator' and e.entity_id = w.entity_id))) or exists (select 1 from w where w.entity_type = 'country' and e.country_iso2 = w.entity_id) or exists (select 1 from w where w.entity_type = 'project' and (e.project_id = w.entity_id or (e.entity_type = 'project' and e.entity_id = w.entity_id))) or exists (select 1 from w where w.entity_type = 'facility' and e.entity_type = 'facility' and e.entity_id = w.entity_id) or exists (select 1 from w where w.entity_type = 'metro' and (e.metro_id = w.entity_id or (e.entity_type = 'facility' and e.entity_id in (select id from facilities where metro_id = w.entity_id)))) ) order by coalesce(e.cluster_id, e.id), e.significance desc, e.detected_at asc ) select ${eventCols(sql)}, count(*) over() as total from events e ${eventJoins(sql)} where e.id in (select id from matched) order by e.detected_at desc, e.id desc limit ${pg_.perPage} offset ${pg_.offset}`; const total = rows.length ? Number(rows[0]!.total) : 0; return { items: await toEventDtos(rows), total, page: pg_.page, perPage: pg_.perPage }; }