diff --git a/ansible/deploy_file_1c_analytics.yml b/ansible/deploy_file_1c_analytics.yml index 4320b9f..dee474e 100644 --- a/ansible/deploy_file_1c_analytics.yml +++ b/ansible/deploy_file_1c_analytics.yml @@ -115,6 +115,7 @@ - check_ingest_freshness.sh - run_company_intelligence_api.sh - run_company_intelligence_refresh.sh + - run_company_registry_bindings_refresh.sh - run_manager_brief.sh - run_ingest_cycle.sh diff --git a/clickhouse-1c/ai/company_intelligence_api.py b/clickhouse-1c/ai/company_intelligence_api.py index 27e462f..66a58d6 100644 --- a/clickhouse-1c/ai/company_intelligence_api.py +++ b/clickhouse-1c/ai/company_intelligence_api.py @@ -463,6 +463,35 @@ def grafana_company_dashboard_url() -> str: ) +def resolve_company_portfolio_card(company_ref: str, infobase: str | None = None) -> dict[str, Any]: + client = ch_client() + filters = [ + "(" + + " OR ".join( + [ + f"company_entity_key = {q(company_ref)}", + f"counterparty = {q(company_ref)}", + f"company_name = {q(company_ref)}", + f"source_counterparty = {q(company_ref)}", + ] + ) + + ")" + ] + if infobase: + filters.append(f"infobase = {q(infobase)}") + sql = f""" + SELECT * + FROM analytics_1c.v_company_portfolio_overview + WHERE {' AND '.join(filters)} + ORDER BY last_company_snapshot_at DESC, amount_30d DESC + LIMIT 1 + """ + rows = rows_to_dict(client.query(sql)) + if not rows: + raise HTTPException(status_code=404, detail="company not found in analytics_1c.v_company_portfolio_overview") + return rows[0] + + def fmt_number(value: Any) -> str: if value is None or value == "": return "-" @@ -484,8 +513,8 @@ def severity_badge(severity: str) -> str: return f'{html.escape(severity or "none")}' -def company_detail_url(counterparty: str, infobase: str | None = None) -> str: - base = f"/manager/company/{quote(counterparty)}" +def company_detail_url(company_ref: str, infobase: str | None = None) -> str: + base = f"/manager/company/{quote(company_ref)}" if infobase: return f"{base}?infobase={quote(infobase)}" return base @@ -538,7 +567,7 @@ def render_manager_brief_html(payload: dict[str, Any]) -> str: f"

{html.escape(item.get('company', '-'))}

{severity_badge(item.get('severity', ''))}
" f"

{html.escape(item.get('reason', '-'))}

" f"

Действие: {html.escape(item.get('recommended_action', '-'))}

" - f"

Открыть карточку компании

" + f"

Открыть карточку компании

" "" ) @@ -550,7 +579,7 @@ def render_manager_brief_html(payload: dict[str, Any]) -> str: f"

{html.escape(company)}

" f"

Прогноз 30д: {html.escape(item.get('forecast_30d', '-'))}

" f"

{html.escape(item.get('interpretation', '-'))}

" - f"Карточка компании" + f"Карточка компании" "" ) @@ -955,7 +984,7 @@ def problematic_companies(days: int = 7, limit: int = 50) -> list[dict[str, Any] WITH recent AS ( SELECT infobase, - counterparty, + counterparty AS company_entity_key, max(generated_at) AS latest_signal_at, max(score) AS max_score, sum(score) AS total_score, @@ -971,7 +1000,9 @@ def problematic_companies(days: int = 7, limit: int = 50) -> list[dict[str, Any] ) SELECT p.infobase AS infobase, + p.company_entity_key AS company_entity_key, p.counterparty AS counterparty, + p.source_counterparty, p.company_name, p.normalized_counterparty, p.registry_match_mode, @@ -997,7 +1028,7 @@ def problematic_companies(days: int = 7, limit: int = 50) -> list[dict[str, Any] FROM recent AS r INNER JOIN analytics_1c.v_company_portfolio_overview AS p ON p.infobase = r.infobase - AND p.counterparty = r.counterparty + AND p.company_entity_key = r.company_entity_key ORDER BY r.max_score DESC, r.signals_total DESC, p.amount_30d DESC, p.counterparty LIMIT {int(limit)} """ @@ -1089,7 +1120,7 @@ def render_problematic_companies_html(items: list[dict[str, Any]], days: int) -> for item in items: rows.append( "" - f"{html.escape(str(item.get('counterparty', '-')))}" + f"{html.escape(str(item.get('counterparty', '-')))}" f"{html.escape(str(item.get('normalized_counterparty') or '-'))}" f"{severity_badge(str(item.get('top_severity') or item.get('signal_severity') or 'none'))}" f"{fmt_number(item.get('max_score'))}" @@ -1231,7 +1262,7 @@ def render_brief_delta_html(payload: dict[str, Any]) -> str: priority_tier = str(item.get("priority_tier") or "low") rows.append( "" - f"{html.escape(str(counterparty or '-'))}" + f"{html.escape(str(counterparty or '-'))}" f"{html.escape(tier_labels.get(priority_tier, priority_tier))}" f"{delta_value(item.get('priority_score', 0))}" f"{html.escape(str(item.get('change_type') or '-'))}" @@ -1538,7 +1569,7 @@ def render_weekly_digest_html(payload: dict[str, Any]) -> str: f"

{html.escape(company)}

{severity_badge(str(item.get('priority') or 'low'))}
" f"

{html.escape(str(item.get('reason') or '-'))}

" f"

Действие: {html.escape(str(item.get('recommended_action') or '-'))}

" - f"

Карточка компании

" + f"

Карточка компании

" "" ) @@ -1792,11 +1823,16 @@ def companies_overview( sql = f""" SELECT infobase, + company_entity_key, organization, counterparty, + source_counterparty, company_name, normalized_counterparty, registry_match_mode, + registry_company_key, + registry_binding_source, + registry_binding_note, registry_assignee_name, registry_status, registry_share_text, @@ -1833,39 +1869,27 @@ def companies_overview( @app.get("/api/1/analytics-1c/companies/{counterparty}/summary") def company_summary(counterparty: str, infobase: str | None = None) -> dict[str, Any]: client = ch_client() - filters = [f"counterparty = {q(counterparty)}"] - if infobase: - filters.append(f"infobase = {q(infobase)}") - sql = f""" - SELECT * - FROM analytics_1c.v_company_portfolio_overview - WHERE {' AND '.join(filters)} - ORDER BY last_company_snapshot_at DESC, amount_30d DESC - LIMIT 1 - """ - rows = rows_to_dict(client.query(sql)) - if not rows: - raise HTTPException(status_code=404, detail="counterparty not found in analytics_1c.v_company_portfolio_overview") - card = rows[0] + card = resolve_company_portfolio_card(counterparty, infobase) + entity_key = str(card.get("company_entity_key") or "") forecast_sql = f""" SELECT metric, horizon_days, baseline_daily, trend_slope, predicted_daily, predicted_total, confidence, note FROM analytics_1c.v_company_forecasts_current - WHERE counterparty = {q(counterparty)} - {"AND infobase = " + q(infobase) if infobase else ""} + WHERE counterparty = {q(entity_key)} + {"AND infobase = " + q(str(card.get('infobase') or infobase)) if (card.get('infobase') or infobase) else ""} ORDER BY metric, horizon_days """ signals_sql = f""" SELECT generated_at, severity, score, signal_type, summary FROM analytics_1c.v_company_health_current - WHERE counterparty = {q(counterparty)} - {"AND infobase = " + q(infobase) if infobase else ""} + WHERE counterparty = {q(entity_key)} + {"AND infobase = " + q(str(card.get('infobase') or infobase)) if (card.get('infobase') or infobase) else ""} ORDER BY score DESC, generated_at DESC """ timeline_sql = f""" SELECT last_company_snapshot_at AS ts, infobase, company_name, owner_user, current_status, db_size_bytes, reglog_size_bytes, active_locks, current_activity_score FROM analytics_1c.v_company_portfolio_overview - WHERE counterparty = {q(counterparty)} - {"AND infobase = " + q(infobase) if infobase else ""} + WHERE company_entity_key = {q(entity_key)} + {"AND infobase = " + q(str(card.get('infobase') or infobase)) if (card.get('infobase') or infobase) else ""} ORDER BY ts DESC LIMIT 1 """ @@ -1875,14 +1899,14 @@ def company_summary(counterparty: str, infobase: str | None = None) -> dict[str, timeline_sql = f""" SELECT ts, infobase, doc_type, operation_type, amount, status, author FROM analytics_1c.documents - WHERE counterparty = {q(counterparty)} - {"AND infobase = " + q(infobase) if infobase else ""} + WHERE infobase = {q(str(card.get('infobase') or infobase or ''))} + AND counterparty != '' ORDER BY ts DESC LIMIT 20 """ timeline = rows_to_dict(client.query(timeline_sql)) essence = ( - f"Компания {counterparty}: за 30 дней событий {card['docs_30d']}, суммарная активность {card['amount_30d']}, " + f"Компания {card['counterparty']}: за 30 дней событий {card['docs_30d']}, суммарная активность {card['amount_30d']}, " f"прогноз активности на 30 дней {card['amount_forecast_30d']}, риск {card['signal_severity']}." ) payload = { @@ -1905,9 +1929,10 @@ def company_forecast( horizon_days: int | None = Query(default=None, ge=1, le=365), ) -> dict[str, Any]: client = ch_client() - filters = [f"counterparty = {q(counterparty)}"] - if infobase: - filters.append(f"infobase = {q(infobase)}") + card = resolve_company_portfolio_card(counterparty, infobase) + filters = [f"counterparty = {q(str(card.get('company_entity_key') or counterparty))}"] + if card.get("infobase") or infobase: + filters.append(f"infobase = {q(str(card.get('infobase') or infobase))}") if horizon_days is not None: filters.append(f"horizon_days = {int(horizon_days)}") sql = f""" @@ -1927,9 +1952,8 @@ def company_timeline( limit: int = Query(default=100, ge=1, le=500), ) -> dict[str, Any]: client = ch_client() - filters = [f"counterparty = {q(counterparty)}"] - if infobase: - filters.append(f"infobase = {q(infobase)}") + card = resolve_company_portfolio_card(counterparty, infobase) + filters = [f"infobase = {q(str(card.get('infobase') or infobase or ''))}", "counterparty != ''"] sql = f""" SELECT ts, infobase, organization, doc_type, doc_number, author, operation_type, amount, status, posted FROM analytics_1c.documents @@ -2076,13 +2100,14 @@ def render_company_detail_html(summary_payload: dict[str, Any], infobase: str | title = card.get("counterparty", "Карточка компании") subtitle = summary_payload.get("essence", "") grafana_url = grafana_company_dashboard_url() - summary_url = f"/api/1/analytics-1c/companies/{quote(card['counterparty'])}/summary" + company_ref = str(card.get("company_entity_key") or card.get("counterparty") or "") + summary_url = f"/api/1/analytics-1c/companies/{quote(company_ref)}/summary" if infobase: summary_url += f"?infobase={quote(infobase)}" - timeline_url = f"/api/1/analytics-1c/companies/{quote(card['counterparty'])}/timeline" + timeline_url = f"/api/1/analytics-1c/companies/{quote(company_ref)}/timeline" if infobase: timeline_url += f"?infobase={quote(infobase)}" - forecast_url = f"/api/1/analytics-1c/companies/{quote(card['counterparty'])}/forecast" + forecast_url = f"/api/1/analytics-1c/companies/{quote(company_ref)}/forecast" if infobase: forecast_url += f"?infobase={quote(infobase)}" diff --git a/clickhouse-1c/ai/generate_manager_brief.py b/clickhouse-1c/ai/generate_manager_brief.py index b25cfd8..c76e443 100644 --- a/clickhouse-1c/ai/generate_manager_brief.py +++ b/clickhouse-1c/ai/generate_manager_brief.py @@ -185,14 +185,14 @@ def snapshot_from_context(context: dict[str, Any]) -> dict[tuple[Any, Any], dict snapshot_items = context.get("portfolio_snapshot") or [] if snapshot_items: return { - (item.get("infobase"), item.get("counterparty")): item + (item.get("infobase"), item.get("company_entity_key") or item.get("counterparty")): item for item in snapshot_items } merged: dict[tuple[Any, Any], dict[str, Any]] = {} for source_name in ("top_risks", "top_forecasts", "watchlist", "busy_bases"): for item in context.get(source_name, []): - key = (item.get("infobase"), item.get("counterparty")) + key = (item.get("infobase"), item.get("company_entity_key") or item.get("counterparty")) if key not in merged: merged[key] = dict(item) else: @@ -268,6 +268,7 @@ def build_context(client, top_limit: int, freshness_hours: int) -> dict[str, Any f""" SELECT infobase, + company_entity_key, counterparty, normalized_counterparty, registry_match_mode, @@ -294,6 +295,7 @@ def build_context(client, top_limit: int, freshness_hours: int) -> dict[str, Any f""" SELECT infobase, + company_entity_key, counterparty, normalized_counterparty, registry_match_mode, @@ -316,6 +318,7 @@ def build_context(client, top_limit: int, freshness_hours: int) -> dict[str, Any f""" SELECT s.infobase, + p.company_entity_key, s.counterparty, p.normalized_counterparty, p.registry_match_mode, @@ -328,7 +331,7 @@ def build_context(client, top_limit: int, freshness_hours: int) -> dict[str, Any round(p.amount_forecast_30d, 2) AS amount_forecast_30d FROM analytics_1c.v_company_health_current AS s LEFT JOIN analytics_1c.v_company_portfolio_overview AS p - ON p.infobase = s.infobase AND p.counterparty = s.counterparty + ON p.infobase = s.infobase AND p.company_entity_key = s.counterparty WHERE s.signal_type IN ('inactive_company', 'amount_drop', 'docs_stopped') ORDER BY s.score DESC, p.days_since_last_activity DESC, p.amount_30d DESC LIMIT {int(top_limit)} @@ -341,6 +344,7 @@ def build_context(client, top_limit: int, freshness_hours: int) -> dict[str, Any f""" SELECT infobase, + company_entity_key, counterparty, normalized_counterparty, current_status, @@ -364,11 +368,14 @@ def build_context(client, top_limit: int, freshness_hours: int) -> dict[str, Any SELECT c.opened_at, c.infobase, + p.company_entity_key, c.entity_id AS counterparty, c.title, c.severity, c.status FROM analytics_1c.cases AS c + LEFT JOIN analytics_1c.v_company_portfolio_overview AS p + ON p.infobase = c.infobase AND p.company_entity_key = c.entity_id WHERE c.entity_type = 'counterparty' AND c.status != 'closed' ORDER BY c.opened_at DESC LIMIT {int(top_limit)} @@ -381,6 +388,7 @@ def build_context(client, top_limit: int, freshness_hours: int) -> dict[str, Any """ SELECT infobase, + company_entity_key, counterparty, normalized_counterparty, registry_match_mode, @@ -424,8 +432,8 @@ def compute_delta_context(current: dict[str, Any], previous_artifact: dict[str, previous = previous_artifact.get("context", {}) current_summary = current.get("portfolio_summary", {}) previous_summary = previous.get("portfolio_summary", {}) - current_watchlist = {(item.get("infobase"), item.get("counterparty")) for item in current.get("watchlist", [])} - previous_watchlist = {(item.get("infobase"), item.get("counterparty")) for item in previous.get("watchlist", [])} + current_watchlist = {(item.get("infobase"), item.get("company_entity_key") or item.get("counterparty")) for item in current.get("watchlist", [])} + previous_watchlist = {(item.get("infobase"), item.get("company_entity_key") or item.get("counterparty")) for item in previous.get("watchlist", [])} current_snapshot = snapshot_from_context(current) previous_snapshot = snapshot_from_context(previous) @@ -530,6 +538,7 @@ def compute_delta_context(current: dict[str, Any], previous_artifact: dict[str, rank_top_change( { "infobase": current_item.get("infobase"), + "company_entity_key": current_item.get("company_entity_key"), "company": current_item.get("counterparty"), "normalized_counterparty": current_item.get("normalized_counterparty"), "registry_match_mode": current_item.get("registry_match_mode"), diff --git a/clickhouse-1c/ai/generate_recovery_brief.py b/clickhouse-1c/ai/generate_recovery_brief.py index 214dddc..6b86381 100644 --- a/clickhouse-1c/ai/generate_recovery_brief.py +++ b/clickhouse-1c/ai/generate_recovery_brief.py @@ -146,7 +146,9 @@ def build_context(client, args: argparse.Namespace) -> dict[str, Any]: ) SELECT p.infobase AS infobase, + p.company_entity_key AS company_entity_key, p.counterparty AS counterparty, + p.source_counterparty, p.company_name, p.normalized_counterparty, p.registry_match_mode, @@ -220,7 +222,7 @@ def render_deterministic_recovery(context: dict[str, Any]) -> dict[str, Any]: ] top_incidents = [] for item in problematic[:6]: - company = str(item.get("counterparty") or "-") + company = str(item.get("counterparty") or item.get("company_name") or "-") actions = [ "Проверить владельца и состав открытых кейсов по компании.", "Подтвердить, что по компании есть план снижения хвоста в ближайшие 24 часа.", @@ -231,6 +233,7 @@ def render_deterministic_recovery(context: dict[str, Any]) -> dict[str, Any]: actions.append("Сначала подтвердить корректность manual-сопоставления.") top_incidents.append( { + "company_entity_key": str(item.get("company_entity_key") or ""), "company": company, "severity": str(item.get("signal_severity") or item.get("top_severity") or "critical"), "diagnosis": ( diff --git a/clickhouse-1c/ai/refresh_company_intelligence.py b/clickhouse-1c/ai/refresh_company_intelligence.py index 14d2de3..45b6831 100644 --- a/clickhouse-1c/ai/refresh_company_intelligence.py +++ b/clickhouse-1c/ai/refresh_company_intelligence.py @@ -130,18 +130,18 @@ def main() -> int: daily_rows = query_rows( client, """ - SELECT infobase, organization, counterparty, d, docs_total, amount_total - FROM analytics_1c.v_counterparty_daily - ORDER BY infobase, counterparty, d + SELECT infobase, organization, company_entity_key, source_counterparty, d, docs_total, amount_total + FROM analytics_1c.v_company_activity_daily + ORDER BY infobase, company_entity_key, d """, ) if not daily_rows: - print("no counterparty rows in analytics_1c.v_counterparty_daily; nothing to refresh") + print("no company activity rows in analytics_1c.v_company_activity_daily; nothing to refresh") return 0 - grouped: dict[tuple[str, str, str], list[DailyPoint]] = defaultdict(list) + grouped: dict[tuple[str, str, str, str], list[DailyPoint]] = defaultdict(list) for row in daily_rows: - key = (row["infobase"], row["organization"], row["counterparty"]) + key = (row["infobase"], row["organization"], row["company_entity_key"], row.get("source_counterparty") or row["company_entity_key"]) grouped[key].append( DailyPoint( d=row["d"], @@ -151,26 +151,26 @@ def main() -> int: ) cases_map = { - (row["infobase"], row["counterparty"]): int(row["open_cases_total"] or 0) + (row["infobase"], row["company_entity_key"]): int(row["open_cases_total"] or 0) for row in query_rows( client, """ - SELECT infobase, entity_id AS counterparty, countIf(status != 'closed') AS open_cases_total + SELECT infobase, entity_id AS company_entity_key, countIf(status != 'closed') AS open_cases_total FROM analytics_1c.cases WHERE entity_type = 'counterparty' - GROUP BY infobase, counterparty + GROUP BY infobase, company_entity_key """, ) } detections_map = { - (row["infobase"], row["counterparty"]): int(row["detections_total"] or 0) + (row["infobase"], row["company_entity_key"]): int(row["detections_total"] or 0) for row in query_rows( client, """ - SELECT infobase, entity_id AS counterparty, count() AS detections_total + SELECT infobase, entity_id AS company_entity_key, count() AS detections_total FROM analytics_1c.detections WHERE entity_type = 'counterparty' AND status != 'closed' - GROUP BY infobase, counterparty + GROUP BY infobase, company_entity_key """, ) } @@ -205,8 +205,8 @@ def main() -> int: forecast_rows: list[list[Any]] = [] signal_rows: list[list[Any]] = [] - for (infobase, _organization, counterparty), points in grouped.items(): - if normalize_company_key(counterparty) in excluded_company_keys: + for (infobase, _organization, company_entity_key, source_counterparty), points in grouped.items(): + if normalize_company_key(source_counterparty) in excluded_company_keys: continue points.sort(key=lambda p: p.d) filled = fill_daily_series(points) @@ -223,8 +223,8 @@ def main() -> int: amount_7d = float(sum(p.amount_total for p in last_7)) amount_prev_7d = float(sum(p.amount_total for p in prev_7)) days_since_last_activity = (date.today() - latest_day).days - open_cases_total = cases_map.get((infobase, counterparty), 0) - detections_total = detections_map.get((infobase, counterparty), 0) + open_cases_total = cases_map.get((infobase, company_entity_key), 0) + detections_total = detections_map.get((infobase, company_entity_key), 0) company_state = company_state_map.get(infobase, {}) current_status = str(company_state.get("current_status") or "") active_locks = int(company_state.get("active_locks") or 0) @@ -245,7 +245,7 @@ def main() -> int: generated_at, latest_day, infobase, - counterparty, + company_entity_key, int(horizon), metric, float(baseline), @@ -261,29 +261,29 @@ def main() -> int: signals: list[tuple[str, int, str, str]] = [] if days_since_last_activity >= 14 and (docs_prev_7d > 0 or amount_prev_7d > 0): - signals.append(("inactive_company", 85, "high", f"Нет активности по компании {counterparty} уже {days_since_last_activity} дн.")) + signals.append(("inactive_company", 85, "high", f"Нет активности по компании {source_counterparty} уже {days_since_last_activity} дн.")) if amount_prev_7d > 0 and amount_7d < amount_prev_7d * 0.5: - signals.append(("amount_drop", 70, "high", f"Активность по компании {counterparty} упала более чем на 50% неделя к неделе.")) + signals.append(("amount_drop", 70, "high", f"Активность по компании {source_counterparty} упала более чем на 50% неделя к неделе.")) if docs_prev_7d > 0 and docs_7d == 0: - signals.append(("docs_stopped", 55, "medium", f"По компании {counterparty} прекратился поток документов за последние 7 дней.")) + signals.append(("docs_stopped", 55, "medium", f"По компании {source_counterparty} прекратился поток документов за последние 7 дней.")) if current_status == "busy" or active_locks > 0 or temp_db_present > 0: score = min(85, 45 + active_locks * 5 + temp_db_present * 10) - signals.append(("base_busy", score, severity_score_to_label(score), f"Файловая база компании {counterparty} занята: status={current_status}, locks={active_locks}, tempDb={temp_db_present}.")) + signals.append(("base_busy", score, severity_score_to_label(score), f"Файловая база компании {source_counterparty} занята: status={current_status}, locks={active_locks}, tempDb={temp_db_present}.")) if scheduler_touched > 0 and current_activity_score >= 15: - signals.append(("scheduler_activity", 35, "medium", f"По компании {counterparty} есть активность scheduler и повышенный activity score {current_activity_score}.")) + signals.append(("scheduler_activity", 35, "medium", f"По компании {source_counterparty} есть активность scheduler и повышенный activity score {current_activity_score}.")) if open_cases_total > 0: - signals.append(("open_cases", min(95, 40 + open_cases_total * 10), severity_score_to_label(min(95, 40 + open_cases_total * 10)), f"По компании {counterparty} есть открытые кейсы: {open_cases_total}.")) + signals.append(("open_cases", min(95, 40 + open_cases_total * 10), severity_score_to_label(min(95, 40 + open_cases_total * 10)), f"По компании {source_counterparty} есть открытые кейсы: {open_cases_total}.")) if detections_total > 0: - signals.append(("open_detections", min(90, 35 + detections_total * 5), severity_score_to_label(min(90, 35 + detections_total * 5)), f"По компании {counterparty} есть активные detections: {detections_total}.")) + signals.append(("open_detections", min(90, 35 + detections_total * 5), severity_score_to_label(min(90, 35 + detections_total * 5)), f"По компании {source_counterparty} есть активные detections: {detections_total}.")) for signal_type, score, severity, summary in signals: signal_rows.append( - [ - generated_at, - infobase, - counterparty, - f"{signal_type}:{infobase}:{counterparty}", - severity, + [ + generated_at, + infobase, + company_entity_key, + f"{signal_type}:{infobase}:{company_entity_key}", + severity, int(score), signal_type, summary, diff --git a/clickhouse-1c/ai/refresh_company_registry_bindings.py b/clickhouse-1c/ai/refresh_company_registry_bindings.py new file mode 100644 index 0000000..12134ea --- /dev/null +++ b/clickhouse-1c/ai/refresh_company_registry_bindings.py @@ -0,0 +1,121 @@ +#!/usr/bin/env python3 +from __future__ import annotations + +import argparse +import os +from datetime import UTC, datetime +from typing import Any + +import clickhouse_connect + + +def parse_args() -> argparse.Namespace: + p = argparse.ArgumentParser(description="Refresh technical company->registry bindings for analytics_1c") + p.add_argument("--host", default=os.getenv("CLICKHOUSE_HOST", "localhost")) + p.add_argument("--port", type=int, default=int(os.getenv("CLICKHOUSE_PORT", "8123"))) + p.add_argument("--user", default=os.getenv("CLICKHOUSE_USER", "default")) + p.add_argument("--password", default=os.getenv("CLICKHOUSE_PASSWORD", "")) + p.add_argument("--database", default=os.getenv("CLICKHOUSE_DB", "analytics_1c")) + return p.parse_args() + + +def ch_client(args: argparse.Namespace): + return clickhouse_connect.get_client( + host=args.host, + port=args.port, + username=args.user, + password=args.password, + database=args.database, + ) + + +def query_rows(client, sql: str) -> list[dict[str, Any]]: + result = client.query(sql) + return [dict(zip(result.column_names, row)) for row in result.result_rows] + + +def main() -> int: + args = parse_args() + client = ch_client(args) + generated_at = datetime.now(UTC).replace(tzinfo=None, microsecond=0) + + current_bindings = { + str(row["company_entity_key"]): row + for row in query_rows( + client, + """ + SELECT + company_entity_key, + registry_company_key, + registry_company_name, + binding_source + FROM analytics_1c.v_company_registry_bindings_current + """, + ) + } + + candidates = query_rows( + client, + """ + SELECT + company_entity_key, + infobase, + base_id, + base_path, + ifNull(base_path_key, '') AS base_path_key, + registry_company_key, + company_name, + registry_match_mode + FROM analytics_1c.v_company_portfolio_overview + WHERE registry_match_mode IN ('direct', 'alias', 'manual') + AND registry_company_key != '' + AND company_entity_key != '' + """ + ) + + inserts: list[list[Any]] = [] + for row in candidates: + entity_key = str(row["company_entity_key"]) + registry_key = str(row["registry_company_key"]) + current = current_bindings.get(entity_key) + if current and str(current.get("registry_company_key") or "") == registry_key: + continue + inserts.append( + [ + generated_at, + str(row["infobase"] or ""), + entity_key, + str(row["base_id"] or ""), + str(row["base_path"] or ""), + str(row["base_path_key"] or ""), + registry_key, + str(row["company_name"] or ""), + f"bootstrap_{row['registry_match_mode']}", + "autobound_from_portfolio", + ] + ) + + if inserts: + client.insert( + "analytics_1c.company_registry_bindings", + inserts, + column_names=[ + "ts", + "infobase", + "company_entity_key", + "base_id", + "base_path", + "base_path_key", + "registry_company_key", + "registry_company_name", + "binding_source", + "note", + ], + ) + + print(f"company registry bindings refreshed: inserted={len(inserts)} generated_at={generated_at.isoformat()}") + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/clickhouse-1c/clickhouse/init/04_company_intelligence.sql b/clickhouse-1c/clickhouse/init/04_company_intelligence.sql index 2aa13b7..3534829 100644 --- a/clickhouse-1c/clickhouse/init/04_company_intelligence.sql +++ b/clickhouse-1c/clickhouse/init/04_company_intelligence.sql @@ -39,38 +39,101 @@ CREATE TABLE IF NOT EXISTS analytics_1c.company_health_signals ENGINE = MergeTree ORDER BY (generated_at, severity, infobase, counterparty, signal_id); -CREATE OR REPLACE VIEW analytics_1c.v_counterparty_daily AS -SELECT - toDate(ts) AS d, - infobase, - organization, - counterparty, - count() AS docs_total, - sum(amount) AS amount_total, - countIf(posted = 1) AS posted_docs_total, - countIf(posted = 0) AS unposted_docs_total, - countIf(status = 'busy') AS busy_docs_total, - countIf(status = 'online') AS online_docs_total, - uniqExact(doc_type) AS doc_types_total -FROM analytics_1c.documents -WHERE counterparty != '' -GROUP BY d, infobase, organization, counterparty; +CREATE TABLE IF NOT EXISTS analytics_1c.company_registry_bindings +( + ts DateTime, + infobase LowCardinality(String), + company_entity_key String, + base_id String, + base_path String, + base_path_key String, + registry_company_key String, + registry_company_name String, + binding_source LowCardinality(String), + note String +) +ENGINE = MergeTree +ORDER BY (company_entity_key, ts); -CREATE OR REPLACE VIEW analytics_1c.v_counterparty_latest_activity AS +CREATE OR REPLACE VIEW analytics_1c.v_companies_current AS SELECT infobase, + company_name, organization, - counterparty, - max(ts) AS last_seen_at, - argMax(doc_type, ts) AS last_doc_type, - argMax(operation_type, ts) AS last_operation_type, - argMax(status, ts) AS last_status, - argMax(amount, ts) AS last_amount, + owner_user, + base_id, + base_path, + trimBoth(replaceRegexpAll(replaceRegexpAll(replaceRegexpAll(upperUTF8(base_path), '[\\\\/]+', '/'), '[^0-9A-ZА-ЯЁ:/._ -]+', ' '), '\\s+', ' ')) AS base_path_key, + multiIf( + base_id != '', concat('baseid:', base_id), + base_path != '', concat('basepath:', trimBoth(replaceRegexpAll(replaceRegexpAll(replaceRegexpAll(upperUTF8(base_path), '[\\\\/]+', '/'), '[^0-9A-ZА-ЯЁ:/._ -]+', ' '), '\\s+', ' '))), + concat('infobase:', trimBoth(replaceRegexpAll(replaceRegexpAll(replaceRegexpAll(upperUTF8(infobase), '(^|\\s)20[0-9]{2}($|\\s)', ' '), '[^0-9A-ZА-ЯЁ]+', ' '), '\\s+', ' '))) + ) AS company_entity_key, + current_status, + db_size_bytes, + reglog_size_bytes, + active_locks, + temp_db_present, + scheduler_touched, + current_activity_score, + last_company_snapshot_at +FROM +( + SELECT + infobase, + argMax(company_name, ts) AS company_name, + argMax(organization, ts) AS organization, + argMax(owner_user, ts) AS owner_user, + argMax(base_id, ts) AS base_id, + argMax(base_path, ts) AS base_path, + argMax(status, ts) AS current_status, + argMax(db_size_bytes, ts) AS db_size_bytes, + argMax(reglog_size_bytes, ts) AS reglog_size_bytes, + argMax(active_locks, ts) AS active_locks, + argMax(temp_db_present, ts) AS temp_db_present, + argMax(scheduler_touched, ts) AS scheduler_touched, + argMax(activity_score, ts) AS current_activity_score, + max(ts) AS last_company_snapshot_at + FROM analytics_1c.companies + GROUP BY infobase +); + +CREATE OR REPLACE VIEW analytics_1c.v_company_activity_daily AS +SELECT + toDate(documents.ts) AS d, + documents.infobase AS infobase, + ifNull(companies.organization, documents.organization) AS organization, + ifNull(companies.company_entity_key, concat('infobase:', trimBoth(replaceRegexpAll(replaceRegexpAll(replaceRegexpAll(upperUTF8(documents.infobase), '(^|\\s)20[0-9]{2}($|\\s)', ' '), '[^0-9A-ZА-ЯЁ]+', ' '), '\\s+', ' ')))) AS company_entity_key, + argMax(documents.counterparty, documents.ts) AS source_counterparty, + count() AS docs_total, + sum(documents.amount) AS amount_total, + countIf(documents.posted = 1) AS posted_docs_total, + countIf(documents.posted = 0) AS unposted_docs_total, + countIf(documents.status = 'busy') AS busy_docs_total, + countIf(documents.status = 'online') AS online_docs_total, + uniqExact(documents.doc_type) AS doc_types_total +FROM analytics_1c.documents AS documents +LEFT JOIN analytics_1c.v_companies_current AS companies ON companies.infobase = documents.infobase +WHERE documents.counterparty != '' +GROUP BY d, documents.infobase, ifNull(companies.organization, documents.organization), ifNull(companies.company_entity_key, concat('infobase:', trimBoth(replaceRegexpAll(replaceRegexpAll(replaceRegexpAll(upperUTF8(documents.infobase), '(^|\\s)20[0-9]{2}($|\\s)', ' '), '[^0-9A-ZА-ЯЁ]+', ' '), '\\s+', ' ')))); + +CREATE OR REPLACE VIEW analytics_1c.v_company_activity_latest AS +SELECT + documents.infobase AS infobase, + ifNull(companies.organization, documents.organization) AS organization, + ifNull(companies.company_entity_key, concat('infobase:', trimBoth(replaceRegexpAll(replaceRegexpAll(replaceRegexpAll(upperUTF8(documents.infobase), '(^|\\s)20[0-9]{2}($|\\s)', ' '), '[^0-9A-ZА-ЯЁ]+', ' '), '\\s+', ' ')))) AS company_entity_key, + argMax(documents.counterparty, documents.ts) AS source_counterparty, + max(documents.ts) AS last_seen_at, + argMax(documents.doc_type, documents.ts) AS last_doc_type, + argMax(documents.operation_type, documents.ts) AS last_operation_type, + argMax(documents.status, documents.ts) AS last_status, + argMax(documents.amount, documents.ts) AS last_amount, count() AS docs_lifetime, - sum(amount) AS amount_lifetime -FROM analytics_1c.documents -WHERE counterparty != '' -GROUP BY infobase, organization, counterparty; + sum(documents.amount) AS amount_lifetime +FROM analytics_1c.documents AS documents +LEFT JOIN analytics_1c.v_companies_current AS companies ON companies.infobase = documents.infobase +WHERE documents.counterparty != '' +GROUP BY documents.infobase, ifNull(companies.organization, documents.organization), ifNull(companies.company_entity_key, concat('infobase:', trimBoth(replaceRegexpAll(replaceRegexpAll(replaceRegexpAll(upperUTF8(documents.infobase), '(^|\\s)20[0-9]{2}($|\\s)', ' '), '[^0-9A-ZА-ЯЁ]+', ' '), '\\s+', ' ')))); CREATE OR REPLACE VIEW analytics_1c.v_company_forecasts_current AS SELECT * @@ -82,24 +145,35 @@ SELECT * FROM analytics_1c.company_health_signals WHERE generated_at = (SELECT max(generated_at) FROM analytics_1c.company_health_signals); -CREATE OR REPLACE VIEW analytics_1c.v_companies_current AS +CREATE OR REPLACE VIEW analytics_1c.v_counterparty_daily AS +SELECT + d, + infobase, + organization, + company_entity_key AS counterparty, + docs_total, + amount_total, + posted_docs_total, + unposted_docs_total, + busy_docs_total, + online_docs_total, + doc_types_total +FROM analytics_1c.v_company_activity_daily; + +CREATE OR REPLACE VIEW analytics_1c.v_counterparty_latest_activity AS SELECT infobase, - argMax(company_name, ts) AS company_name, - argMax(organization, ts) AS organization, - argMax(owner_user, ts) AS owner_user, - argMax(base_id, ts) AS base_id, - argMax(base_path, ts) AS base_path, - argMax(status, ts) AS current_status, - argMax(db_size_bytes, ts) AS db_size_bytes, - argMax(reglog_size_bytes, ts) AS reglog_size_bytes, - argMax(active_locks, ts) AS active_locks, - argMax(temp_db_present, ts) AS temp_db_present, - argMax(scheduler_touched, ts) AS scheduler_touched, - argMax(activity_score, ts) AS current_activity_score, - max(ts) AS last_company_snapshot_at -FROM analytics_1c.companies -GROUP BY infobase; + organization, + company_entity_key AS counterparty, + source_counterparty, + last_seen_at, + last_doc_type, + last_operation_type, + last_status, + last_amount, + docs_lifetime, + amount_lifetime +FROM analytics_1c.v_company_activity_latest; CREATE OR REPLACE VIEW analytics_1c.v_company_registry_current AS SELECT @@ -115,6 +189,21 @@ SELECT FROM analytics_1c.company_registry GROUP BY company_key; +CREATE OR REPLACE VIEW analytics_1c.v_company_registry_bindings_current AS +SELECT + company_entity_key, + argMax(infobase, ts) AS infobase, + argMax(base_id, ts) AS base_id, + argMax(base_path, ts) AS base_path, + argMax(base_path_key, ts) AS base_path_key, + argMax(registry_company_key, ts) AS registry_company_key, + argMax(registry_company_name, ts) AS registry_company_name, + argMax(binding_source, ts) AS binding_source, + argMax(note, ts) AS note, + max(ts) AS last_binding_at +FROM analytics_1c.company_registry_bindings +GROUP BY company_entity_key; + CREATE OR REPLACE VIEW analytics_1c.v_company_registry_alias_map AS SELECT source_company_key, @@ -220,7 +309,7 @@ base AS ( SELECT *, - trimBoth(replaceRegexpAll(replaceRegexpAll(replaceRegexpAll(upperUTF8(counterparty), '(^|\\s)20[0-9]{2}($|\\s)', ' '), '[^0-9A-ZА-ЯЁ]+', ' '), '\\s+', ' ')) AS counterparty_key + trimBoth(replaceRegexpAll(replaceRegexpAll(replaceRegexpAll(upperUTF8(source_counterparty), '(^|\\s)20[0-9]{2}($|\\s)', ' '), '[^0-9A-ZА-ЯЁ]+', ' '), '\\s+', ' ')) AS source_counterparty_key FROM analytics_1c.v_counterparty_latest_activity ), d7 AS @@ -252,6 +341,11 @@ company_state AS SELECT * FROM analytics_1c.v_companies_current ), +binding_state AS +( + SELECT * + FROM analytics_1c.v_company_registry_bindings_current +), registry_state AS ( SELECT * @@ -271,18 +365,18 @@ signals AS ( SELECT infobase, - counterparty, + counterparty AS company_entity_key, max(score) AS signal_score, argMax(severity, score) AS signal_severity, argMax(summary, score) AS top_signal FROM analytics_1c.v_company_health_current - GROUP BY infobase, counterparty + GROUP BY infobase, company_entity_key ), amount_forecast AS ( SELECT infobase, - counterparty, + counterparty AS company_entity_key, predicted_total AS amount_forecast_30d, confidence AS amount_forecast_confidence FROM analytics_1c.v_company_forecasts_current @@ -293,7 +387,7 @@ docs_forecast AS ( SELECT infobase, - counterparty, + counterparty AS company_entity_key, predicted_total AS docs_forecast_30d, confidence AS docs_forecast_confidence FROM analytics_1c.v_company_forecasts_current @@ -304,30 +398,35 @@ cases_current AS ( SELECT infobase, - entity_id AS counterparty, + entity_id AS company_entity_key, countIf(status != 'closed') AS open_cases_total FROM analytics_1c.cases WHERE entity_type = 'counterparty' - GROUP BY infobase, counterparty + GROUP BY infobase, company_entity_key ), detections_current AS ( SELECT infobase, - entity_id AS counterparty, + entity_id AS company_entity_key, count() AS detections_total FROM analytics_1c.detections WHERE entity_type = 'counterparty' AND status != 'closed' - GROUP BY infobase, counterparty + GROUP BY infobase, company_entity_key ) SELECT base.infobase AS infobase, + base.counterparty AS company_entity_key, if(company_state.organization != '', company_state.organization, base.organization) AS organization, - base.counterparty AS counterparty, - if(alias_state.target_company_name != '', alias_state.target_company_name, if(manual_state.company_name != '', manual_state.company_name, if(company_state.company_name != '', company_state.company_name, base.counterparty))) AS company_name, - if(alias_state.target_company_name != '', alias_state.target_company_name, if(manual_state.company_name != '', manual_state.company_name, base.counterparty)) AS normalized_counterparty, - multiIf(ifNull(alias_state.exclude_from_portfolio, 0) = 1, 'excluded', alias_state.target_company_key != '' AND registry_state.company_key != '', 'alias', registry_state.company_key != '', 'direct', manual_state.company_key != '', 'manual', 'none') AS registry_match_mode, + base.source_counterparty AS source_counterparty, + if(binding_state.registry_company_name != '', binding_state.registry_company_name, if(alias_state.target_company_name != '', alias_state.target_company_name, if(manual_state.company_name != '', manual_state.company_name, if(company_state.company_name != '', company_state.company_name, base.source_counterparty)))) AS counterparty, + if(binding_state.registry_company_name != '', binding_state.registry_company_name, if(alias_state.target_company_name != '', alias_state.target_company_name, if(manual_state.company_name != '', manual_state.company_name, if(company_state.company_name != '', company_state.company_name, base.source_counterparty)))) AS company_name, + if(binding_state.registry_company_name != '', binding_state.registry_company_name, if(alias_state.target_company_name != '', alias_state.target_company_name, if(manual_state.company_name != '', manual_state.company_name, base.source_counterparty))) AS normalized_counterparty, + multiIf(binding_state.registry_company_key != '', 'technical', ifNull(alias_state.exclude_from_portfolio, 0) = 1, 'excluded', alias_state.target_company_key != '' AND registry_state.company_key != '', 'alias', registry_state.company_key != '', 'direct', manual_state.company_key != '', 'manual', 'none') AS registry_match_mode, + if(binding_state.registry_company_key != '', binding_state.registry_company_key, if(registry_state.company_key != '', registry_state.company_key, ifNull(manual_state.company_key, ''))) AS registry_company_key, + ifNull(binding_state.binding_source, '') AS registry_binding_source, + ifNull(binding_state.note, '') AS registry_binding_note, if(registry_state.assignee_name != '', registry_state.assignee_name, ifNull(manual_state.assignee_name, '')) AS registry_assignee_name, if(registry_state.registry_status != '', registry_state.registry_status, ifNull(manual_state.registry_status, '')) AS registry_status, if(registry_state.share_text != '', registry_state.share_text, ifNull(manual_state.share_text, '')) AS registry_share_text, @@ -337,6 +436,7 @@ SELECT ifNull(company_state.owner_user, '') AS owner_user, ifNull(company_state.base_id, '') AS base_id, ifNull(company_state.base_path, '') AS base_path, + ifNull(company_state.base_path_key, '') AS base_path_key, base.last_seen_at, company_state.last_company_snapshot_at, base.last_doc_type, @@ -367,14 +467,15 @@ SELECT ifNull(signals.top_signal, '') AS top_signal FROM base LEFT JOIN company_state ON company_state.infobase = base.infobase -LEFT JOIN alias_state ON alias_state.source_company_key = base.counterparty_key -LEFT JOIN registry_state ON registry_state.company_key = if(alias_state.target_company_key != '', alias_state.target_company_key, base.counterparty_key) -LEFT JOIN manual_state ON manual_state.company_key = if(alias_state.target_company_key != '', alias_state.target_company_key, base.counterparty_key) +LEFT JOIN binding_state ON binding_state.company_entity_key = base.counterparty +LEFT JOIN alias_state ON alias_state.source_company_key = base.source_counterparty_key +LEFT JOIN registry_state ON registry_state.company_key = if(binding_state.registry_company_key != '', binding_state.registry_company_key, if(alias_state.target_company_key != '', alias_state.target_company_key, base.source_counterparty_key)) +LEFT JOIN manual_state ON manual_state.company_key = if(binding_state.registry_company_key != '', binding_state.registry_company_key, if(alias_state.target_company_key != '', alias_state.target_company_key, base.source_counterparty_key)) LEFT JOIN d7 ON d7.infobase = base.infobase AND d7.counterparty = base.counterparty LEFT JOIN d30 ON d30.infobase = base.infobase AND d30.counterparty = base.counterparty -LEFT JOIN signals ON signals.infobase = base.infobase AND signals.counterparty = base.counterparty -LEFT JOIN amount_forecast ON amount_forecast.infobase = base.infobase AND amount_forecast.counterparty = base.counterparty -LEFT JOIN docs_forecast ON docs_forecast.infobase = base.infobase AND docs_forecast.counterparty = base.counterparty -LEFT JOIN cases_current ON cases_current.infobase = base.infobase AND cases_current.counterparty = base.counterparty -LEFT JOIN detections_current ON detections_current.infobase = base.infobase AND detections_current.counterparty = base.counterparty +LEFT JOIN signals ON signals.infobase = base.infobase AND signals.company_entity_key = base.counterparty +LEFT JOIN amount_forecast ON amount_forecast.infobase = base.infobase AND amount_forecast.company_entity_key = base.counterparty +LEFT JOIN docs_forecast ON docs_forecast.infobase = base.infobase AND docs_forecast.company_entity_key = base.counterparty +LEFT JOIN cases_current ON cases_current.infobase = base.infobase AND cases_current.company_entity_key = base.counterparty +LEFT JOIN detections_current ON detections_current.infobase = base.infobase AND detections_current.company_entity_key = base.counterparty WHERE ifNull(alias_state.exclude_from_portfolio, 0) = 0; diff --git a/clickhouse-1c/detections/build_entity_timeline.sql b/clickhouse-1c/detections/build_entity_timeline.sql index 7ddc597..50833fc 100644 --- a/clickhouse-1c/detections/build_entity_timeline.sql +++ b/clickhouse-1c/detections/build_entity_timeline.sql @@ -21,19 +21,22 @@ INSERT INTO analytics_1c.entity_timeline SELECT * FROM ( SELECT - ts, + documents.ts AS ts, 'counterparty' AS entity_type, - counterparty AS entity_id, - infobase, - author AS actor, + ifNull(portfolio.company_entity_key, documents.counterparty) AS entity_id, + documents.infobase AS infobase, + documents.author AS actor, 'documents' AS source, - concat('counterparty:', operation_type) AS event_type, - if(status = 'busy', 'medium', 'low') AS severity, - greatest(10, toUInt32(round(amount))) AS score, - concat('counterparty:', counterparty, ':', doc_id, ':', toString(toUnixTimestamp(ts))) AS ref_id, - concat('Активность компании ', counterparty, ': ', doc_type, ' score=', toString(amount), ' status=', status) AS summary - FROM analytics_1c.documents - WHERE counterparty != '' + concat('counterparty:', documents.operation_type) AS event_type, + if(documents.status = 'busy', 'medium', 'low') AS severity, + greatest(10, toUInt32(round(documents.amount))) AS score, + concat('counterparty:', ifNull(portfolio.company_entity_key, documents.counterparty), ':', documents.doc_id, ':', toString(toUnixTimestamp(documents.ts))) AS ref_id, + concat('Активность компании ', ifNull(portfolio.company_name, documents.counterparty), ': ', documents.doc_type, ' score=', toString(documents.amount), ' status=', documents.status) AS summary + FROM analytics_1c.documents AS documents + LEFT JOIN analytics_1c.v_company_portfolio_overview AS portfolio + ON portfolio.infobase = documents.infobase + AND portfolio.source_counterparty = documents.counterparty + WHERE documents.counterparty != '' ) AS src WHERE src.ref_id NOT IN (SELECT ref_id FROM analytics_1c.entity_timeline); @@ -41,18 +44,18 @@ INSERT INTO analytics_1c.entity_timeline SELECT * FROM ( SELECT - ts, + last_company_snapshot_at AS ts, 'counterparty' AS entity_type, - company_name AS entity_id, + company_entity_key AS entity_id, infobase, owner_user AS actor, 'companies' AS source, 'company_snapshot' AS event_type, - if(status = 'busy' OR active_locks > 0 OR temp_db_present = 1, 'medium', 'low') AS severity, - greatest(10, toUInt32(round(activity_score))) AS score, - concat('company:', infobase, ':', toString(toUnixTimestamp(ts))) AS ref_id, - concat('Company snapshot ', company_name, ': status=', status, ' locks=', toString(active_locks), ' score=', toString(activity_score)) AS summary - FROM analytics_1c.companies + if(current_status = 'busy' OR active_locks > 0 OR temp_db_present = 1, 'medium', 'low') AS severity, + greatest(10, toUInt32(round(current_activity_score))) AS score, + concat('company:', company_entity_key, ':', toString(toUnixTimestamp(last_company_snapshot_at))) AS ref_id, + concat('Company snapshot ', company_name, ': status=', current_status, ' locks=', toString(active_locks), ' score=', toString(current_activity_score)) AS summary + FROM analytics_1c.v_companies_current ) AS src WHERE src.ref_id NOT IN (SELECT ref_id FROM analytics_1c.entity_timeline); diff --git a/clickhouse-1c/ops/run_company_intelligence_refresh.sh b/clickhouse-1c/ops/run_company_intelligence_refresh.sh index d3b9651..19930de 100644 --- a/clickhouse-1c/ops/run_company_intelligence_refresh.sh +++ b/clickhouse-1c/ops/run_company_intelligence_refresh.sh @@ -35,6 +35,8 @@ docker exec -i "${CH_CONTAINER}" clickhouse-client \ --database "${CLICKHOUSE_DB}" \ < "${ROOT}/clickhouse/init/04_company_intelligence.sql" +"${ROOT}/ops/run_company_registry_bindings_refresh.sh" + "${VENV}/bin/python" "${ROOT}/ai/refresh_company_intelligence.py" \ --host "${CH_RUNTIME_HOST}" \ --port "${CLICKHOUSE_PORT}" \ diff --git a/clickhouse-1c/ops/run_company_registry_bindings_refresh.sh b/clickhouse-1c/ops/run_company_registry_bindings_refresh.sh new file mode 100644 index 0000000..9cefeeb --- /dev/null +++ b/clickhouse-1c/ops/run_company_registry_bindings_refresh.sh @@ -0,0 +1,31 @@ +#!/usr/bin/env bash +set -euo pipefail + +ROOT="${AW_1C_ROOT:-/opt/activitywatch/clickhouse-1c}" +ENV_FILE="${ROOT}/.env" +VENV="${ROOT}/.venv" + +if [[ ! -f "${ENV_FILE}" ]]; then + echo "missing env file: ${ENV_FILE}" >&2 + exit 1 +fi + +if [[ ! -x "${VENV}/bin/python" ]]; then + echo "missing venv python: ${VENV}/bin/python" >&2 + exit 1 +fi + +# shellcheck disable=SC1090 +. "${ENV_FILE}" + +CH_RUNTIME_HOST="${AW_1C_CLICKHOUSE_RUNTIME_HOST:-${CLICKHOUSE_HOST}}" +if [[ "${CH_RUNTIME_HOST}" == "clickhouse" ]]; then + CH_RUNTIME_HOST="127.0.0.1" +fi + +"${VENV}/bin/python" "${ROOT}/ai/refresh_company_registry_bindings.py" \ + --host "${CH_RUNTIME_HOST}" \ + --port "${CLICKHOUSE_PORT}" \ + --user "${CLICKHOUSE_USER}" \ + --password "${CLICKHOUSE_PASSWORD}" \ + --database "${CLICKHOUSE_DB}" diff --git a/clickhouse-1c/ops/run_ingest_cycle.sh b/clickhouse-1c/ops/run_ingest_cycle.sh index 3926359..72669b7 100644 --- a/clickhouse-1c/ops/run_ingest_cycle.sh +++ b/clickhouse-1c/ops/run_ingest_cycle.sh @@ -55,6 +55,7 @@ docker exec -i "${CH_CONTAINER}" clickhouse-client \ --database "${CLICKHOUSE_DB}" \ < "${ROOT}/clickhouse/init/04_company_intelligence.sql" +"${ROOT}/ops/run_company_registry_bindings_refresh.sh" "${ROOT}/ops/run_company_intelligence_refresh.sh" docker exec -i "${CH_CONTAINER}" clickhouse-client \