diff --git a/clickhouse-1c/README.md b/clickhouse-1c/README.md index 05e6b73..dc5ceb8 100644 --- a/clickhouse-1c/README.md +++ b/clickhouse-1c/README.md @@ -26,6 +26,8 @@ File 1C + reglog + host telemetry ├─ raw_* ├─ documents ├─ postings + ├─ business_events + ├─ document_change_events ├─ reglog_events ├─ audit_events ├─ host_events @@ -45,6 +47,8 @@ File 1C + reglog + host telemetry - `clickhouse/init/*.sql` — схема БД. - `etl/load_1c_exports.py` — loader CSV/JSON выгрузок в raw/core таблицы. - `etl/config.example.yml` — пример ETL-конфига. +- `docs/1C_BUSINESS_EVENT_LAYER_RU.md` — production contract следующего шага: + canonical business-event слой для документов/проводок/изменений. - `detections/rules.yml` — каталог правил detections. - `detections/insert_detections.sql` — SQL-шаблоны rule-based detections. - `grafana/dashboard-catalog.md` — целевая структура дашбордов. @@ -95,7 +99,7 @@ docker compose up -d 3. Инициализировать landing-каталоги и ETL config: ```bash -mkdir -p landing/{documents,postings,reglog,audit,host} +mkdir -p landing/{documents,postings,business_events,document_changes,companies,reglog,audit,host} cp etl/config.example.yml etl/config.yml ``` @@ -156,6 +160,8 @@ Read-only API для руководителя: - выгрузки 1С по документам; - выгрузки движений/проводок; +- выгрузки canonical business events; +- выгрузки изменений документов и реквизитов; - read-only `companies` snapshot по файловым базам; - журнал регистрации 1С; - audit/export критичных изменений; @@ -184,3 +190,14 @@ Read-only API для руководителя: - case/timeline слой считается вне 1С. - это прогноз активности компании/базы, а не финансовых проводок; - если `counterparty` в live-выгрузках пустой, company-forecast слой останется корректно пустым. + +## Следующий production шаг + +В репо уже заложен scaffold под следующий слой: + +- `business_events` — единый event stream бухгалтерских событий; +- `document_change_events` — изменения документов/реквизитов; +- ETL уже умеет принимать эти datasets в `landing/business_events` и + `landing/document_changes`; +- дальше нужен только read-only extractor, который будет наполнять их из 1С + или внешних безопасных выгрузок без записи в базу. diff --git a/clickhouse-1c/clickhouse/init/01_raw_tables.sql b/clickhouse-1c/clickhouse/init/01_raw_tables.sql index 81ac911..9e528eb 100644 --- a/clickhouse-1c/clickhouse/init/01_raw_tables.sql +++ b/clickhouse-1c/clickhouse/init/01_raw_tables.sql @@ -16,6 +16,24 @@ CREATE TABLE IF NOT EXISTS analytics_1c.raw_1c_postings ENGINE = MergeTree ORDER BY (ingested_at, source_file); +CREATE TABLE IF NOT EXISTS analytics_1c.raw_1c_business_events +( + ingested_at DateTime DEFAULT now(), + source_file String, + payload String +) +ENGINE = MergeTree +ORDER BY (ingested_at, source_file); + +CREATE TABLE IF NOT EXISTS analytics_1c.raw_1c_document_changes +( + ingested_at DateTime DEFAULT now(), + source_file String, + payload String +) +ENGINE = MergeTree +ORDER BY (ingested_at, source_file); + CREATE TABLE IF NOT EXISTS analytics_1c.raw_1c_companies ( ingested_at DateTime DEFAULT now(), diff --git a/clickhouse-1c/clickhouse/init/02_core_tables.sql b/clickhouse-1c/clickhouse/init/02_core_tables.sql index 670b5f9..8d68751 100644 --- a/clickhouse-1c/clickhouse/init/02_core_tables.sql +++ b/clickhouse-1c/clickhouse/init/02_core_tables.sql @@ -32,6 +32,56 @@ CREATE TABLE IF NOT EXISTS analytics_1c.postings ENGINE = MergeTree ORDER BY (infobase, ts, registrar); +CREATE TABLE IF NOT EXISTS analytics_1c.business_events +( + ts DateTime, + event_id String, + infobase LowCardinality(String), + company_entity_key String, + organization String, + department String, + document_id String, + document_number String, + document_type LowCardinality(String), + registrar String, + operation_type String, + event_kind LowCardinality(String), + user String, + counterparty String, + counterparty_inn String, + debit_account String, + credit_account String, + amount Decimal(18, 2), + currency LowCardinality(String), + line_no UInt32, + evidence_ref String, + source_file String +) +ENGINE = MergeTree +ORDER BY (infobase, ts, document_id, line_no, event_id); + +CREATE TABLE IF NOT EXISTS analytics_1c.document_change_events +( + ts DateTime, + change_id String, + infobase LowCardinality(String), + company_entity_key String, + organization String, + document_id String, + document_number String, + document_type LowCardinality(String), + change_kind LowCardinality(String), + field_name String, + user String, + before_value String, + after_value String, + risk_tag String, + evidence_ref String, + source_file String +) +ENGINE = MergeTree +ORDER BY (infobase, ts, document_id, change_id); + CREATE TABLE IF NOT EXISTS analytics_1c.companies ( ts DateTime, diff --git a/clickhouse-1c/etl/config.example.yml b/clickhouse-1c/etl/config.example.yml index b5f459d..25817f0 100644 --- a/clickhouse-1c/etl/config.example.yml +++ b/clickhouse-1c/etl/config.example.yml @@ -8,6 +8,8 @@ clickhouse: landing: documents: ./landing/documents postings: ./landing/postings + business_events: ./landing/business_events + document_changes: ./landing/document_changes companies: ./landing/companies reglog: ./landing/reglog audit: ./landing/audit @@ -17,6 +19,8 @@ formats: default: jsonl documents: jsonl postings: jsonl + business_events: jsonl + document_changes: jsonl companies: jsonl reglog: jsonl audit: jsonl diff --git a/clickhouse-1c/etl/load_1c_exports.py b/clickhouse-1c/etl/load_1c_exports.py index ce19062..287215e 100644 --- a/clickhouse-1c/etl/load_1c_exports.py +++ b/clickhouse-1c/etl/load_1c_exports.py @@ -17,6 +17,8 @@ from dateutil import parser as date_parser RAW_TABLES = { "documents": "raw_1c_documents", "postings": "raw_1c_postings", + "business_events": "raw_1c_business_events", + "document_changes": "raw_1c_document_changes", "companies": "raw_1c_companies", "reglog": "raw_reglog", "audit": "raw_audit", @@ -26,6 +28,8 @@ RAW_TABLES = { CORE_TABLES = { "documents": "documents", "postings": "postings", + "business_events": "business_events", + "document_changes": "document_change_events", "companies": "companies", "reglog": "reglog_events", "audit": "audit_events", @@ -46,7 +50,7 @@ class Config: def parse_args() -> argparse.Namespace: p = argparse.ArgumentParser(description="Load file-based 1C exports into ClickHouse") p.add_argument("--config", required=True, help="Path to YAML config") - p.add_argument("--dataset", choices=["documents", "postings", "companies", "reglog", "audit", "host"], help="Load only one dataset") + p.add_argument("--dataset", choices=["documents", "postings", "business_events", "document_changes", "companies", "reglog", "audit", "host"], help="Load only one dataset") return p.parse_args() @@ -129,6 +133,50 @@ def map_core_row(dataset: str, source_file: str, row: dict[str, Any]) -> list[An float(row.get("amount", 0) or 0), source_file, ] + if dataset == "business_events": + return [ + normalize_ts(row.get("ts") or row.get("event_time")), + row.get("event_id", ""), + row.get("infobase", ""), + row.get("company_entity_key", ""), + row.get("organization", ""), + row.get("department", ""), + row.get("document_id", row.get("doc_id", "")), + row.get("document_number", row.get("doc_number", "")), + row.get("document_type", row.get("doc_type", "")), + row.get("registrar", ""), + row.get("operation_type", ""), + row.get("event_kind", ""), + row.get("user", row.get("author", "")), + row.get("counterparty", ""), + row.get("counterparty_inn", ""), + row.get("debit_account", row.get("account_dt", "")), + row.get("credit_account", row.get("account_ct", "")), + float(row.get("amount", 0) or 0), + row.get("currency", "RUB"), + int(row.get("line_no", 0) or 0), + row.get("evidence_ref", ""), + source_file, + ] + if dataset == "document_changes": + return [ + normalize_ts(row.get("ts") or row.get("change_time")), + row.get("change_id", ""), + row.get("infobase", ""), + row.get("company_entity_key", ""), + row.get("organization", ""), + row.get("document_id", row.get("doc_id", "")), + row.get("document_number", row.get("doc_number", "")), + row.get("document_type", row.get("doc_type", "")), + row.get("change_kind", ""), + row.get("field_name", ""), + row.get("user", row.get("author", "")), + row.get("before_value", ""), + row.get("after_value", ""), + row.get("risk_tag", ""), + row.get("evidence_ref", ""), + source_file, + ] if dataset == "companies": return [ normalize_ts(row.get("ts")), @@ -194,6 +242,10 @@ def core_columns(dataset: str) -> list[str]: return ["ts", "infobase", "organization", "department", "doc_type", "doc_id", "doc_number", "author", "counterparty", "operation_type", "amount", "status", "posted", "source_file"] if dataset == "postings": return ["ts", "infobase", "registrar", "operation_type", "account_dt", "account_ct", "amount", "source_file"] + if dataset == "business_events": + return ["ts", "event_id", "infobase", "company_entity_key", "organization", "department", "document_id", "document_number", "document_type", "registrar", "operation_type", "event_kind", "user", "counterparty", "counterparty_inn", "debit_account", "credit_account", "amount", "currency", "line_no", "evidence_ref", "source_file"] + if dataset == "document_changes": + return ["ts", "change_id", "infobase", "company_entity_key", "organization", "document_id", "document_number", "document_type", "change_kind", "field_name", "user", "before_value", "after_value", "risk_tag", "evidence_ref", "source_file"] if dataset == "companies": return ["ts", "infobase", "company_name", "organization", "owner_user", "base_id", "base_path", "status", "db_size_bytes", "reglog_size_bytes", "active_locks", "temp_db_present", "scheduler_touched", "activity_score", "source_file"] if dataset == "reglog": diff --git a/clickhouse-1c/etl/test_load_1c_exports.py b/clickhouse-1c/etl/test_load_1c_exports.py new file mode 100644 index 0000000..f3f58e6 --- /dev/null +++ b/clickhouse-1c/etl/test_load_1c_exports.py @@ -0,0 +1,81 @@ +#!/usr/bin/env python3 +from __future__ import annotations + +import sys +import unittest +from pathlib import Path +from types import SimpleNamespace + +sys.path.insert(0, str(Path(__file__).resolve().parent)) +sys.modules.setdefault("clickhouse_connect", SimpleNamespace(get_client=lambda **_: None)) + +import load_1c_exports as loader + + +class Load1CExportsTests(unittest.TestCase): + def test_business_events_mapping(self) -> None: + row = { + "event_time": "2026-05-22T12:00:00Z", + "event_id": "evt-1", + "infobase": "ИБ-1", + "company_entity_key": "baseid:abc", + "organization": "ООО Тест", + "department": "Продажи", + "doc_id": "doc-1", + "doc_number": "0001", + "doc_type": "Реализация", + "registrar": "DOC-1", + "operation_type": "sale", + "event_kind": "posting", + "author": "user1", + "counterparty": "Контрагент", + "counterparty_inn": "7700000000", + "account_dt": "62.01", + "account_ct": "90.01", + "amount": "125000.50", + "line_no": "2", + "evidence_ref": "reglog:1", + } + mapped = loader.map_core_row("business_events", "events.jsonl", row) + self.assertEqual(mapped[1], "evt-1") + self.assertEqual(mapped[3], "baseid:abc") + self.assertEqual(mapped[7], "0001") + self.assertEqual(mapped[15], "62.01") + self.assertEqual(mapped[17], 125000.5) + self.assertEqual(mapped[19], 2) + self.assertEqual(mapped[21], "events.jsonl") + + def test_document_changes_mapping(self) -> None: + row = { + "change_time": "2026-05-22T12:00:00Z", + "change_id": "chg-1", + "infobase": "ИБ-1", + "company_entity_key": "basepath:c:/1c/base", + "organization": "ООО Тест", + "document_id": "doc-2", + "document_number": "0002", + "document_type": "Поступление", + "change_kind": "requisites_change", + "field_name": "Контрагент", + "user": "user2", + "before_value": "Старый", + "after_value": "Новый", + "risk_tag": "counterparty_change", + "evidence_ref": "audit:2", + } + mapped = loader.map_core_row("document_changes", "changes.jsonl", row) + self.assertEqual(mapped[1], "chg-1") + self.assertEqual(mapped[3], "basepath:c:/1c/base") + self.assertEqual(mapped[8], "requisites_change") + self.assertEqual(mapped[12], "Новый") + self.assertEqual(mapped[15], "changes.jsonl") + + def test_new_dataset_columns_shape(self) -> None: + self.assertEqual(loader.RAW_TABLES["business_events"], "raw_1c_business_events") + self.assertEqual(loader.CORE_TABLES["document_changes"], "document_change_events") + self.assertEqual(len(loader.core_columns("business_events")), 22) + self.assertEqual(len(loader.core_columns("document_changes")), 16) + + +if __name__ == "__main__": + unittest.main() diff --git a/clickhouse-1c/ops/bootstrap_runtime.sh b/clickhouse-1c/ops/bootstrap_runtime.sh index 1123a07..74e69ff 100644 --- a/clickhouse-1c/ops/bootstrap_runtime.sh +++ b/clickhouse-1c/ops/bootstrap_runtime.sh @@ -6,6 +6,8 @@ ROOT="${AW_1C_ROOT:-/opt/activitywatch/clickhouse-1c}" mkdir -p \ "${ROOT}/landing/documents" \ "${ROOT}/landing/postings" \ + "${ROOT}/landing/business_events" \ + "${ROOT}/landing/document_changes" \ "${ROOT}/landing/companies" \ "${ROOT}/landing/registry" \ "${ROOT}/landing/reglog" \ @@ -13,6 +15,8 @@ mkdir -p \ "${ROOT}/landing/host" \ "${ROOT}/archive/documents" \ "${ROOT}/archive/postings" \ + "${ROOT}/archive/business_events" \ + "${ROOT}/archive/document_changes" \ "${ROOT}/archive/companies" \ "${ROOT}/archive/registry" \ "${ROOT}/archive/reglog" \ diff --git a/docs/1C_BUSINESS_EVENT_LAYER_RU.md b/docs/1C_BUSINESS_EVENT_LAYER_RU.md new file mode 100644 index 0000000..73598c0 --- /dev/null +++ b/docs/1C_BUSINESS_EVENT_LAYER_RU.md @@ -0,0 +1,192 @@ +# 1C Business Event Layer + +Это следующий production-шаг поверх уже работающего контура +`file 1C -> ClickHouse -> Grafana -> AI Investigator`. + +Цель не в том, чтобы заменить текущий `company intelligence`, а в том, чтобы +добавить **read-only business-event слой**, пригодный для финансовых +расследований, explainability и rule-based detections по бухгалтерскому смыслу. + +## Почему этот слой нужен + +Текущий контур уже решает: + +- operational telemetry; +- timeline по reglog/audit; +- detections/cases; +- company portfolio intelligence; +- manager/recovery briefs. + +Но он ещё не даёт полноценного ответа на вопросы уровня: + +- какие проводки дали вклад в аномалию; +- кто и когда перепровёл документ; +- какие изменения реквизитов повлияли на результат; +- почему вырос НДС, возвраты или нетипичные движения. + +Для этого нужен отдельный канонический слой бизнес-событий. + +## Слои данных + +### 1. `documents` + +Карточки документов и базовые агрегаты по документам. + +### 2. `postings` + +Лёгкий слой проводок/движений. Уже есть в контуре, но он недостаточен как +канонический event stream. + +### 3. `business_events` + +Новый целевой канонический слой. Один ряд = одно бизнес-событие, пригодное для: + +- timeline; +- detections; +- explainability; +- AI investigations; +- correlation с reglog/audit/cases. + +Текущая schema scaffold: + +- `ts` +- `event_id` +- `infobase` +- `company_entity_key` +- `organization` +- `department` +- `document_id` +- `document_number` +- `document_type` +- `registrar` +- `operation_type` +- `event_kind` +- `user` +- `counterparty` +- `counterparty_inn` +- `debit_account` +- `credit_account` +- `amount` +- `currency` +- `line_no` +- `evidence_ref` +- `source_file` + +### 4. `document_change_events` + +Новый слой изменений документов и реквизитов. + +Нужен для: + +- расследования перепроведений; +- контроля изменений реквизитов; +- reconstruction narrative; +- объяснения, какие именно изменения дали бизнес-эффект. + +Текущая schema scaffold: + +- `ts` +- `change_id` +- `infobase` +- `company_entity_key` +- `organization` +- `document_id` +- `document_number` +- `document_type` +- `change_kind` +- `field_name` +- `user` +- `before_value` +- `after_value` +- `risk_tag` +- `evidence_ref` +- `source_file` + +## Read-only extraction path + +Правильный extraction path такой: + +1. Внешний extractor читает только безопасные read-only источники. +2. Формирует `jsonl/csv` в landing-каталоги. +3. `etl/load_1c_exports.py` грузит данные в raw/core ClickHouse tables. +4. Detection/AI слой работает только с ClickHouse. + +То есть LLM и manager pages не ходят в 1С напрямую. + +## Безопасные источники для extractor + +Подходящие: + +- регламентированные выгрузки документов/движений; +- журнал регистрации 1С; +- внешние реестры и справочники; +- audit/export критичных изменений; +- отдельные read-only файлы, формируемые рядом с 1С. + +Не подходящие по умолчанию: + +- запись в файловую базу; +- опасные `COM`/`Configurator` сценарии; +- любой write-back path в production 1С. + +## Почему здесь нужен `company_entity_key` + +Этот слой не должен зависеть от того, как компания названа в текущий момент в +1С. Поэтому новые event tables уже сразу завязаны на `company_entity_key`. + +Приоритет идентификации: + +1. `baseid:` +2. `basepath:` +3. только fallback на human name + +Это делает timeline и расследования устойчивыми к rename. + +## Что extractor должен уметь первым + +Минимальный read-only extractor v1 должен уметь: + +- выгружать документы; +- выгружать проводки; +- выгружать `business_events`; +- выгружать `document_change_events`; +- стабильно наполнять `company_entity_key`; +- писать `evidence_ref`, чтобы расследование не было бездоказательным. + +## Первые детекты на этом слое + +После появления business-event выгрузок стоит вводить: + +- ночные проводки; +- дробление платежей; +- повторные перепроведения; +- изменение реквизитов перед/после движения; +- возвраты после закрытия периода; +- нехарактерные движения по пользователю; +- циклические движения по контрагенту/счётам. + +## Что уже сделано в репо + +Уже добавлены: + +- raw tables: + - `raw_1c_business_events` + - `raw_1c_document_changes` +- core tables: + - `business_events` + - `document_change_events` +- ETL support: + - `landing/business_events` + - `landing/document_changes` + - dataset mapping в `etl/load_1c_exports.py` + +Это именно scaffold, а не обещание, что live extractor уже существует. + +## Что делать дальше + +1. Реализовать read-only extractor в отдельном модуле. +2. Стабильно наполнять `company_entity_key`. +3. Добавить первые SQL detections на `business_events`. +4. Обогащать `entity_timeline` уже не только telemetry/audit, но и + business-event evidence. +5. Включить AI Investigator narrative по реальным проводкам и изменениям.