import { createClient, type ClickHouseClient } from "@clickhouse/client"; /** * ClickHouse: append-only analytics store — crawl history, extracted observations, page changes, * connector metrics and daily time series. Postgres remains the transactional source of truth. */ let _ch: ClickHouseClient | null = null; export function getClickHouse(): ClickHouseClient { if (!_ch) { _ch = createClient({ url: process.env.CLICKHOUSE_URL ?? "http://127.0.0.1:8123", database: process.env.CLICKHOUSE_DB ?? "dci", username: process.env.CLICKHOUSE_USER ?? "default", password: process.env.CLICKHOUSE_PASSWORD ?? "", request_timeout: 30_000, clickhouse_settings: { async_insert: 1, wait_for_async_insert: 0 }, }); } return _ch; } export const CLICKHOUSE_DDL: string[] = [ `CREATE TABLE IF NOT EXISTS crawl_log ( ts DateTime64(3) DEFAULT now64(3), connector_id LowCardinality(String), source_id LowCardinality(String), document_id String, url String, fetch_level UInt8, fetcher LowCardinality(String), status_code UInt16, duration_ms UInt32, bytes UInt32, changed UInt8, not_modified UInt8, error LowCardinality(String), credits Float32, run_id String ) ENGINE = MergeTree ORDER BY (connector_id, ts) TTL toDateTime(ts) + INTERVAL 400 DAY`, `CREATE TABLE IF NOT EXISTS observations ( ts DateTime64(3) DEFAULT now64(3), connector_id LowCardinality(String), source_id LowCardinality(String), document_id String, entity_type LowCardinality(String), entity_key String, entity_id String, field LowCardinality(String), value String, value_num Nullable(Float64), confidence LowCardinality(String), is_estimate UInt8, method LowCardinality(String), extractor_version LowCardinality(String), run_id String ) ENGINE = MergeTree ORDER BY (entity_type, entity_key, field, ts)`, `CREATE TABLE IF NOT EXISTS page_changes ( ts DateTime64(3) DEFAULT now64(3), connector_id LowCardinality(String), document_id String, url String, old_hash String, new_hash String, added_lines UInt32, removed_lines UInt32, ratio Float32, significance UInt8, fields Array(String), event_types Array(String) ) ENGINE = MergeTree ORDER BY (connector_id, ts)`, `CREATE TABLE IF NOT EXISTS metrics ( ts DateTime DEFAULT now(), metric LowCardinality(String), dim LowCardinality(String), value Float64 ) ENGINE = MergeTree ORDER BY (metric, dim, ts) TTL ts + INTERVAL 400 DAY`, `CREATE TABLE IF NOT EXISTS entity_daily ( day Date, metric LowCardinality(String), dim LowCardinality(String), value Float64 ) ENGINE = ReplacingMergeTree ORDER BY (metric, dim, day)`, `CREATE TABLE IF NOT EXISTS api_requests ( ts DateTime DEFAULT now(), route LowCardinality(String), status UInt16, duration_ms UInt32, cached UInt8 ) ENGINE = MergeTree ORDER BY (route, ts) TTL ts + INTERVAL 90 DAY`, ]; export async function ensureClickHouse(): Promise { const db = process.env.CLICKHOUSE_DB ?? "dci"; const admin = createClient({ url: process.env.CLICKHOUSE_URL ?? "http://127.0.0.1:8123", username: process.env.CLICKHOUSE_USER ?? "default", password: process.env.CLICKHOUSE_PASSWORD ?? "" }); await admin.command({ query: `CREATE DATABASE IF NOT EXISTS ${db}` }); await admin.close(); const ch = getClickHouse(); for (const ddl of CLICKHOUSE_DDL) await ch.command({ query: ddl }); } export async function chInsert(table: string, rows: Record[]): Promise { if (!rows.length) return; try { await getClickHouse().insert({ table, values: rows, format: "JSONEachRow" }); } catch (e) { if (process.env.DCI_CLICKHOUSE_OPTIONAL !== "0") console.warn(`[clickhouse] insert ${table} failed: ${(e as Error).message}`); else throw e; } } export async function chQuery>(query: string, params: Record = {}): Promise { const rs = await getClickHouse().query({ query, query_params: params, format: "JSONEachRow" }); return (await rs.json()) as T[]; }