From d2415fcfeb9ae630afcf1f04612aeeb0d6cd09c4 Mon Sep 17 00:00:00 2001 From: mi Date: Fri, 21 Aug 2026 11:48:24 +0300 Subject: [PATCH] =?UTF-8?q?=D0=A4=D0=B8=D0=BD=D0=B0=D0=BB=D1=8C=D0=BD?= =?UTF-8?q?=D0=B0=D1=8F=20=D1=81=D1=82=D0=B0=D0=B1=D0=B8=D0=BB=D1=8C=D0=BD?= =?UTF-8?q?=D0=B0=D1=8F=20=D1=80=D0=B5=D0=B0=D0=BB=D0=B8=D0=B7=D0=B0=D1=86?= =?UTF-8?q?=D0=B8=D1=8F?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .gitignore | 5 + VM2_services/codebase/services/.env.example | 1 + .../services/bitrix-sync/.env.example | 1 + .../codebase/services/bitrix-sync/README.md | 9 + .../services/bitrix-sync/app/config.py | 6 + .../services/bitrix-sync/app/engine.py | 176 +++++++++++++++--- .../bitrix-sync/app/reconciliation.py | 43 +++-- .../services/bitrix-sync/app/repository.py | 62 ++++-- .../bitrix-sync/compose.fragment.yaml | 1 + .../services/bitrix-sync/tests/conftest.py | 1 + .../services/bitrix-sync/tests/test_config.py | 5 + .../tests/test_engine_boundaries.py | 113 ++++++++++- .../bitrix-sync/tests/test_repository.py | 123 +++++++++++- .../codebase/services/docker-compose.yml | 2 + .../documentation/module-07-bitrix-sync.md | 1 + VM2_services/vm1-bug-v10.tar.gz | Bin 0 -> 3295 bytes support¬es/backlog.md | 170 +++++++++++++++++ 17 files changed, 662 insertions(+), 57 deletions(-) create mode 100644 .gitignore create mode 100644 VM2_services/vm1-bug-v10.tar.gz diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..03ffb11 --- /dev/null +++ b/.gitignore @@ -0,0 +1,5 @@ +.env +*.log +node_modules/ +__pycache__/ +for_bugs_exchange/ \ No newline at end of file diff --git a/VM2_services/codebase/services/.env.example b/VM2_services/codebase/services/.env.example index 3e8db80..20216a3 100644 --- a/VM2_services/codebase/services/.env.example +++ b/VM2_services/codebase/services/.env.example @@ -30,6 +30,7 @@ BITRIX_SYNC_MODE=disabled BITRIX_SYNC_CONTACT_USER_ID_FIELD=UF_CRM_ BITRIX_SYNC_CONTACT_REGISTERED_FIELD=UF_CRM_ BITRIX_SYNC_CONTACT_CITIZENSHIP_FIELD=UF_CRM_ +BITRIX_SYNC_CONTACT_SOURCE= BITRIX_SYNC_PORTAL_HOST=.bitrix24.ru BITRIX_SYNC_PORTAL_MEMBER_ID= BITRIX_SYNC_PUBLIC_BASE_URL=https:// diff --git a/VM2_services/codebase/services/bitrix-sync/.env.example b/VM2_services/codebase/services/bitrix-sync/.env.example index 03bb32d..0939bab 100644 --- a/VM2_services/codebase/services/bitrix-sync/.env.example +++ b/VM2_services/codebase/services/bitrix-sync/.env.example @@ -11,6 +11,7 @@ BITRIX_SYNC_PUBLIC_BASE_URL=https://sync.example.ru BITRIX_SYNC_CONTACT_USER_ID_FIELD=UF_CRM_100 BITRIX_SYNC_CONTACT_REGISTERED_FIELD=UF_CRM_101 BITRIX_SYNC_CONTACT_CITIZENSHIP_FIELD=UF_CRM_102 +BITRIX_SYNC_CONTACT_SOURCE=WEB BITRIX_SYNC_WEBHOOK_ALLOWED_CIDRS=203.0.113.0/24 BITRIX_SYNC_HTTP_TIMEOUT_SEC=10 BITRIX_SYNC_DB_POOL_SIZE=5 diff --git a/VM2_services/codebase/services/bitrix-sync/README.md b/VM2_services/codebase/services/bitrix-sync/README.md index ae66027..d29fa6e 100644 --- a/VM2_services/codebase/services/bitrix-sync/README.md +++ b/VM2_services/codebase/services/bitrix-sync/README.md @@ -54,6 +54,15 @@ gates: локальный managed PostgreSQL не поднимается Compose 4. Заполнить и активировать валидную `business_alerts` settings version: entity/category/stage/field IDs и responsible party. Placeholder `null` запрещает alert receiver. + +`value_json` настройки `business_alerts` использует ключи `entity_type_id`, +`category_id`, `stage_new`, опциональный `responsible_id` и объект `field_ids`. +В `field_ids` REST-имена пользовательских полей сопоставляются ключам +`alert_number`, `fingerprint`, `alert_type`, `severity`, `app_user_id`, +`selected_external_id`, `occurrence_count`, +`first_occurred_at`, `last_occurred_at`, `workflow_id`. +Конфликтующие Contact передаются в стандартное поле `contactIds`, а значение +`BITRIX_SYNC_CONTACT_SOURCE` — в стандартное поле `sourceId`. 5. Проверить least-privilege credential negative tests; credential администратора запрещён. 6. Валидировать nginx exact routes, no-redirect HTTP policy, body/rate limits, diff --git a/VM2_services/codebase/services/bitrix-sync/app/config.py b/VM2_services/codebase/services/bitrix-sync/app/config.py index 0619bf1..12d9126 100644 --- a/VM2_services/codebase/services/bitrix-sync/app/config.py +++ b/VM2_services/codebase/services/bitrix-sync/app/config.py @@ -46,6 +46,7 @@ class Settings(BaseSettings): contact_user_id_field: str | None = None contact_registered_field: str | None = None contact_citizenship_field: str | None = None + contact_source: str | None = Field(default=None, min_length=1, max_length=128) webhook_allowed_cidrs: str = "" http_timeout_sec: float = Field(default=10, ge=1, le=60) @@ -97,6 +98,7 @@ class Settings(BaseSettings): "contact_user_id_field": self.contact_user_id_field, "contact_registered_field": self.contact_registered_field, "contact_citizenship_field": self.contact_citizenship_field, + "contact_source": self.contact_source, } missing = [name for name, value in required.items() if not value] if missing: @@ -110,6 +112,10 @@ class Settings(BaseSettings): raise ValueError(f"{name} must match UF_CRM_") if not MEMBER_RE.fullmatch(str(self.portal_member_id)): raise ValueError("portal_member_id has invalid format") + assert self.contact_source + self.contact_source = self.contact_source.strip() + if not self.contact_source: + raise ValueError("contact_source cannot be blank") crm = urlsplit(self.crm_rest_webhook_url.get_secret_value()) public = urlsplit(str(self.public_base_url)) diff --git a/VM2_services/codebase/services/bitrix-sync/app/engine.py b/VM2_services/codebase/services/bitrix-sync/app/engine.py index 058f8db..ed7f751 100644 --- a/VM2_services/codebase/services/bitrix-sync/app/engine.py +++ b/VM2_services/codebase/services/bitrix-sync/app/engine.py @@ -15,6 +15,7 @@ from app.domain import ( choose_newest, full_jitter_delay, parse_crm_datetime, + safe_hash, select_email, validate_phone, ) @@ -199,6 +200,14 @@ class WorkflowEngine: alert_code: str | None = None if same_user: selected_id = str(same_user["ID"]) + if len(contacts) > 1: + await self._alert( + workflow_id, + profile.user_id, + "duplicate_contacts", + [str(item["ID"]) for item in contacts], + selected_external_id=selected_id, + ) elif not contacts: selected_id = await self._create_contact(workflow_id, profile) else: @@ -214,15 +223,27 @@ class WorkflowEngine: assert selected is not None all_ids = [item.b24_id for item in candidates] if selected.crm_user_id and selected.crm_user_id != str(profile.user_id): - selected_id = await self._create_contact(workflow_id, profile) alert_code = "contact_owned_by_other_user" + selected_id = await self._create_contact(workflow_id, profile) + await self._alert( + workflow_id, + profile.user_id, + alert_code, + [*all_ids, selected_id], + selected_external_id=selected_id, + ) else: selected_id = selected.b24_id - await self._write_identity(workflow_id, selected_id, profile.user_id, active=True) if len(candidates) > 1: alert_code = "duplicate_contacts" - if alert_code: - await self._alert(workflow_id, profile.user_id, alert_code, all_ids) + await self._alert( + workflow_id, + profile.user_id, + alert_code, + all_ids, + selected_external_id=selected_id, + ) + await self._write_identity(workflow_id, selected_id, profile.user_id, active=True) await self._activate_mapping(workflow_id, profile.user_id, selected_id) async def _create_contact(self, workflow_id: uuid.UUID, profile: Profile) -> str: @@ -551,31 +572,144 @@ class WorkflowEngine: ) async def _alert( - self, workflow_id: uuid.UUID, user_id: uuid.UUID, alert_type: str, candidates: list[str] + self, + workflow_id: uuid.UUID, + user_id: uuid.UUID, + alert_type: str, + candidates: list[str], + *, + selected_external_id: str | None = None, ) -> None: - fingerprint = f"{alert_type}:{user_id}" + fingerprint = safe_hash(f"{alert_type}:{user_id}") + assert fingerprint is not None + async with self.repository.transaction() as connection: + alert = ( + await connection.execute( + text( + """ + INSERT INTO bitrix_sync.business_alerts + (id,fingerprint,alert_type,severity,app_user_id, + selected_external_id,candidate_external_ids,workflow_id,status, + occurrence_count,first_occurred_at,last_occurred_at,created_at,updated_at) + VALUES (gen_random_uuid(),:fingerprint, + :type,'warning',:user_id,:selected_external_id,:candidates, + :workflow_id,'open',1,now(),now(),now(),now()) + ON CONFLICT (alert_type,fingerprint) WHERE status='open' + DO UPDATE SET + selected_external_id=coalesce( + excluded.selected_external_id, + business_alerts.selected_external_id + ), + candidate_external_ids=excluded.candidate_external_ids, + workflow_id=excluded.workflow_id, + occurrence_count=business_alerts.occurrence_count+1, + last_occurred_at=now(),updated_at=now() + RETURNING id,alert_number,fingerprint,alert_type,severity,app_user_id, + selected_external_id,candidate_external_ids,remote_item_id, + occurrence_count,first_occurred_at,last_occurred_at + """ + ), + { + "fingerprint": fingerprint, + "type": alert_type, + "user_id": user_id, + "selected_external_id": selected_external_id, + "candidates": candidates, + "workflow_id": workflow_id, + }, + ) + ).mappings().one() + alert_config = ( + await connection.execute( + text( + """ + SELECT value_json FROM bitrix_sync.settings + WHERE key='business_alerts' AND active=true + AND validation_status='valid' + """ + ) + ) + ).scalar_one_or_none() + + if not isinstance(alert_config, dict): + raise RetryableWorkflow("business_alerts_not_configured") + entity_type_id = alert_config.get("entity_type_id") + stage_new = alert_config.get("stage_new") + if not entity_type_id or not stage_new: + raise RetryableWorkflow("business_alerts_not_configured") + + title = ( + f"Конфликт синхронизации #{alert['alert_number']}: {alert_type}; " + f"user={user_id}; contacts={','.join(candidates)}" + ) + fields: dict[str, Any] = { + "title": title[:255], + "stageId": stage_new, + "contactIds": [int(candidate) for candidate in candidates if candidate.isdigit()], + "sourceId": self.settings.contact_source, + } + if category_id := alert_config.get("category_id"): + fields["categoryId"] = category_id + if responsible_id := alert_config.get("responsible_id"): + fields["assignedById"] = responsible_id + alert_values = { + "alert_number": str(alert["alert_number"]), + "fingerprint": alert["fingerprint"], + "alert_type": alert["alert_type"], + "severity": alert["severity"], + "app_user_id": str(alert["app_user_id"]), + "selected_external_id": alert["selected_external_id"] or "", + "occurrence_count": alert["occurrence_count"], + "first_occurred_at": alert["first_occurred_at"].isoformat(), + "last_occurred_at": alert["last_occurred_at"].isoformat(), + "workflow_id": str(workflow_id), + } + field_ids = alert_config.get("field_ids") + if isinstance(field_ids, dict): + for value_name, field_id in field_ids.items(): + if value_name in alert_values and isinstance(field_id, str) and field_id: + fields[field_id] = alert_values[value_name] + + remote_item_id = alert["remote_item_id"] + if remote_item_id: + await self._command( + workflow_id, + "alert_update", + "crm.item.update", + { + "entityTypeId": int(entity_type_id), + "id": remote_item_id, + "fields": fields, + }, + mutating=True, + ) + return + + result = await self._command( + workflow_id, + "alert_add", + "crm.item.add", + {"entityTypeId": int(entity_type_id), "fields": fields}, + mutating=True, + ) + payload = result.result if isinstance(result.result, dict) else {} + item = payload.get("item", payload) + remote_item_id = item.get("id") if isinstance(item, dict) else None + if remote_item_id is None: + raise RetryableWorkflow("alert_item_id_missing") async with self.repository.transaction() as connection: await connection.execute( text( """ - INSERT INTO bitrix_sync.business_alerts - (id,fingerprint,alert_type,severity,app_user_id,candidate_external_ids, - workflow_id,status,occurrence_count,first_occurred_at,last_occurred_at, - created_at,updated_at) - VALUES (gen_random_uuid(),encode(digest(:fingerprint,'sha256'),'hex'), - :type,'warning',:user_id,:candidates,:workflow_id,'open',1, - now(),now(),now(),now()) - ON CONFLICT (alert_type,fingerprint) WHERE status='open' - DO UPDATE SET occurrence_count=business_alerts.occurrence_count+1, - last_occurred_at=now(),updated_at=now() + UPDATE bitrix_sync.business_alerts + SET remote_item_id=:remote_item_id,remote_stage_id=:stage_id,updated_at=now() + WHERE id=:id """ ), { - "fingerprint": fingerprint, - "type": alert_type, - "user_id": user_id, - "candidates": candidates, - "workflow_id": workflow_id, + "id": alert["id"], + "remote_item_id": str(remote_item_id), + "stage_id": str(stage_new), }, ) diff --git a/VM2_services/codebase/services/bitrix-sync/app/reconciliation.py b/VM2_services/codebase/services/bitrix-sync/app/reconciliation.py index fb71df2..2055cea 100644 --- a/VM2_services/codebase/services/bitrix-sync/app/reconciliation.py +++ b/VM2_services/codebase/services/bitrix-sync/app/reconciliation.py @@ -45,13 +45,13 @@ class IncrementalReconciler: started_at = datetime.now(UTC) since = (cursor or datetime(1970, 1, 1, tzinfo=UTC)) - timedelta(seconds=overlap_seconds) start = 0 - scanned: list[str] = [] + scanned: list[tuple[str, str]] = [] while True: result = await self.crm.call( "crm.item.list", { "entityTypeId": 3, - "select": ["id"], + "select": ["id", "updatedTime"], "filter": { ">=updatedTime": since.replace(tzinfo=None).isoformat(timespec="seconds"), "opened": 1, @@ -63,10 +63,15 @@ class IncrementalReconciler: ) if result.outcome != CrmOutcome.SUCCEEDED: raise RuntimeError(result.error_code or "reconciliation_failed") - payload: dict[str, Any] = result.result or {} - items = payload.get("items", payload if isinstance(payload, list) else []) - scanned.extend(str(item["id"]) for item in items if "id" in item) - next_start = payload.get("next") + payload: Any = result.result or {} + items = payload.get("items", []) if isinstance(payload, dict) else payload + if not isinstance(items, list): + raise RuntimeError("reconciliation_items_invalid") + for item in items: + if not isinstance(item, dict) or "id" not in item or "updatedTime" not in item: + raise RuntimeError("reconciliation_item_missing_identity") + scanned.append((str(item["id"]), str(item["updatedTime"]))) + next_start = payload.get("next") if isinstance(payload, dict) else None if next_start is None: break start = int(next_start) @@ -89,27 +94,35 @@ class IncrementalReconciler: ) return len(scanned) - async def _enqueue_changed(self, external_ids: list[str]) -> None: - if not external_ids: + async def _enqueue_changed(self, changed_contacts: list[tuple[str, str]]) -> None: + if not changed_contacts: return + external_ids = [external_id for external_id, _ in changed_contacts] + event_ids = [ + f"contact.reconciliation:{external_id}:{updated_time}" + for external_id, updated_time in changed_contacts + ] async with self.repository.transaction() as connection: await connection.execute( text( """ INSERT INTO bitrix_sync.webhook_inbox - (id,receiver_type,event_type,external_entity_id,status,coalesced_count, - received_at,last_received_at) - SELECT gen_random_uuid(),'contact','contact.reconciliation',id,'received',1, - now(),now() - FROM unnest(CAST(:ids AS text[])) id + (id,receiver_type,event_type,event_id,external_entity_id,status, + coalesced_count,received_at,last_received_at) + SELECT gen_random_uuid(),'contact','contact.reconciliation',candidate.event_id, + candidate.external_id,'received',1,now(),now() + FROM unnest(CAST(:ids AS text[]),CAST(:event_ids AS text[])) + AS candidate(external_id,event_id) WHERE NOT EXISTS ( SELECT 1 FROM bitrix_sync.webhook_inbox w - WHERE w.receiver_type='contact' AND w.external_entity_id=id + WHERE w.receiver_type='contact' + AND w.external_entity_id=candidate.external_id AND w.status IN ('received','processing','retry_wait') ) + ON CONFLICT (receiver_type,event_id) WHERE event_id IS NOT NULL DO NOTHING """ ), - {"ids": external_ids}, + {"ids": external_ids, "event_ids": event_ids}, ) diff --git a/VM2_services/codebase/services/bitrix-sync/app/repository.py b/VM2_services/codebase/services/bitrix-sync/app/repository.py index a452e9c..a7e8ea8 100644 --- a/VM2_services/codebase/services/bitrix-sync/app/repository.py +++ b/VM2_services/codebase/services/bitrix-sync/app/repository.py @@ -12,6 +12,8 @@ from typing import Any from sqlalchemy import text from sqlalchemy.ext.asyncio import AsyncConnection, AsyncEngine, create_async_engine +from app.domain import safe_hash + def postgres_ssl_context() -> ssl.SSLContext: ca_file = os.environ.get("PG_CA_FILE") @@ -83,9 +85,12 @@ class Repository: WITH candidates AS ( SELECT id FROM han_app.sync_queue - WHERE status IN ('pending','retry_wait') - AND next_attempt_at <= now() - AND (locked_until IS NULL OR locked_until < now()) + WHERE ( + status IN ('pending','retry_wait') AND next_attempt_at <= now() + ) + OR ( + status='leased' AND locked_until < now() + ) ORDER BY next_attempt_at, created_at FOR UPDATE SKIP LOCKED LIMIT :limit @@ -156,7 +161,7 @@ class Repository: """ ) async with self.engine.begin() as connection: - return ( + workflow_id = ( await connection.execute( sql, { @@ -167,6 +172,17 @@ class Repository: }, ) ).scalar_one() + await connection.execute( + text( + """ + UPDATE bitrix_sync.workflow_instances + SET state='running',current_step='load_profile',updated_at=now() + WHERE id=:id AND state IN ('created','running') + """ + ), + {"id": workflow_id}, + ) + return workflow_id async def complete_task(self, task: LeasedTask, workflow_id: uuid.UUID) -> bool: async with self.engine.begin() as connection: @@ -280,8 +296,12 @@ class Repository: """ WITH candidates AS ( SELECT id FROM bitrix_sync.webhook_inbox - WHERE status IN ('received','retry_wait') AND next_attempt_at<=now() - AND (locked_until IS NULL OR locked_until Settings: contact_user_id_field="UF_CRM_100", contact_registered_field="UF_CRM_101", contact_citizenship_field="UF_CRM_102", + contact_source="WEB", webhook_allowed_cidrs="203.0.113.0/24", ) diff --git a/VM2_services/codebase/services/bitrix-sync/tests/test_config.py b/VM2_services/codebase/services/bitrix-sync/tests/test_config.py index 70a7cd9..bf510a8 100644 --- a/VM2_services/codebase/services/bitrix-sync/tests/test_config.py +++ b/VM2_services/codebase/services/bitrix-sync/tests/test_config.py @@ -31,6 +31,7 @@ def test_full_mode_rejects_portal_host_mismatch() -> None: contact_user_id_field="UF_CRM_1", contact_registered_field="UF_CRM_2", contact_citizenship_field="UF_CRM_3", + contact_source="WEB", webhook_allowed_cidrs="203.0.113.0/24", ) @@ -56,6 +57,10 @@ def test_image_default_allows_compose_process_role_override() -> None: for source in (compose, production_compose): assert 'command: ["han-bitrix-sync-worker"]' in source assert 'command: ["han-bitrix-sync-reconciliation"]' in source + reconciliation_service = source.split(" bitrix-sync-reconciliation:", 1)[1].split( + "\n bitrix-sync-", 1 + )[0] + assert 'restart: "no"' in reconciliation_service def test_alembic_chain_preserves_legacy_baseline() -> None: diff --git a/VM2_services/codebase/services/bitrix-sync/tests/test_engine_boundaries.py b/VM2_services/codebase/services/bitrix-sync/tests/test_engine_boundaries.py index 1b98d7a..2bbfa64 100644 --- a/VM2_services/codebase/services/bitrix-sync/tests/test_engine_boundaries.py +++ b/VM2_services/codebase/services/bitrix-sync/tests/test_engine_boundaries.py @@ -2,6 +2,7 @@ from __future__ import annotations import uuid from contextlib import asynccontextmanager +from datetime import UTC, datetime import pytest from sqlalchemy.dialects.postgresql import JSONB @@ -14,9 +15,54 @@ from app.repository import Profile class Result: rowcount = 1 + def __init__(self, *, row=None, scalar=None) -> None: + self.row = row + self.scalar = scalar + + def mappings(self): + return self + + def one(self): + return self.row + + def scalar_one_or_none(self): + return self.scalar + class Connection: + def __init__(self, statements: list[str]) -> None: + self.statements = statements + async def execute(self, statement, params=None): + sql = str(statement) + self.statements.append(sql) + if "INSERT INTO bitrix_sync.business_alerts" in sql: + now = datetime.now(UTC) + return Result( + row={ + "id": uuid.uuid4(), + "alert_number": 1, + "fingerprint": "fingerprint", + "alert_type": params["type"], + "severity": "warning", + "app_user_id": params["user_id"], + "selected_external_id": params["selected_external_id"], + "candidate_external_ids": params["candidates"], + "remote_item_id": None, + "occurrence_count": 1, + "first_occurred_at": now, + "last_occurred_at": now, + } + ) + if "SELECT value_json FROM bitrix_sync.settings" in sql: + return Result( + scalar={ + "entity_type_id": 178, + "category_id": 5, + "stage_new": "DT178_5:NEW", + "field_ids": {"candidate_external_ids": "ufCrm_999"}, + } + ) return Result() @@ -27,7 +73,7 @@ class FakeRepository: @asynccontextmanager async def transaction(self): - yield Connection() + yield Connection(self.statements) async def active_mapping(self, user_id): return self.mapping @@ -37,8 +83,16 @@ class FakeRepository: class FakeCrm: - def __init__(self, user_id: uuid.UUID) -> None: + def __init__( + self, + user_id: uuid.UUID, + *, + existing_identity: bool = False, + foreign_owner: bool = False, + ) -> None: self.user_id = user_id + self.existing_identity = existing_identity + self.foreign_owner = foreign_owner self.calls: list[tuple[str, dict]] = [] async def call(self, method, params, *, mutating): @@ -52,9 +106,19 @@ class FakeCrm: { "ID": contact_id, "CREATED_TIME": "2026-08-06T10:00:00Z", - "UF_CRM_100": None, + "UF_CRM_100": ( + str(self.user_id) + if self.existing_identity and contact_id == "10" + else str(uuid.UUID(int=1)) + if self.foreign_owner and contact_id == "10" + else None + ), }, ) + if method == "crm.contact.add": + return CrmResult(CrmOutcome.SUCCEEDED, "11") + if method == "crm.item.add": + return CrmResult(CrmOutcome.SUCCEEDED, {"item": {"id": "500"}}) return CrmResult(CrmOutcome.SUCCEEDED, True) @@ -76,3 +140,46 @@ async def test_multiple_contacts_choose_numeric_newest(full_settings) -> None: updates = [params for method, params in crm.calls if method == "crm.contact.update"] assert updates[0]["id"] == "10" assert not any(method == "crm.contact.add" for method, _ in crm.calls) + assert any(method == "crm.item.add" for method, _ in crm.calls) + alert_add = next(params for method, params in crm.calls if method == "crm.item.add") + assert alert_add["fields"]["contactIds"] == [9, 10] + assert alert_add["fields"]["sourceId"] == "WEB" + assert "ufCrm_999" not in alert_add["fields"] + assert any("entity_external_mapping" in statement for statement in repository.statements) + assert not any("digest(" in statement for statement in repository.statements) + + +@pytest.mark.asyncio +async def test_recovered_duplicate_creates_alert_before_mapping(full_settings) -> None: + user_id = uuid.uuid4() + repository = FakeRepository() + crm = FakeCrm(user_id, existing_identity=True) + engine = WorkflowEngine(repository, crm, full_settings) + + await engine._map_or_create( + uuid.uuid4(), + Profile(user_id=user_id, phone="+79001234567", identity_status="A", profile_status="A"), + ) + + methods = [method for method, _ in crm.calls] + assert "crm.item.add" in methods + assert "crm.contact.update" not in methods + assert any("entity_external_mapping" in statement for statement in repository.statements) + + +@pytest.mark.asyncio +async def test_foreign_owned_contact_alert_includes_new_contact(full_settings) -> None: + user_id = uuid.uuid4() + repository = FakeRepository() + crm = FakeCrm(user_id, foreign_owner=True) + engine = WorkflowEngine(repository, crm, full_settings) + + await engine._map_or_create( + uuid.uuid4(), + Profile(user_id=user_id, phone="+79001234567", identity_status="A", profile_status="A"), + ) + + methods = [method for method, _ in crm.calls] + alert_add = next(params for method, params in crm.calls if method == "crm.item.add") + assert methods.index("crm.contact.add") < methods.index("crm.item.add") + assert alert_add["fields"]["contactIds"] == [9, 10, 11] diff --git a/VM2_services/codebase/services/bitrix-sync/tests/test_repository.py b/VM2_services/codebase/services/bitrix-sync/tests/test_repository.py index 8e02b49..59c288d 100644 --- a/VM2_services/codebase/services/bitrix-sync/tests/test_repository.py +++ b/VM2_services/codebase/services/bitrix-sync/tests/test_repository.py @@ -1,11 +1,15 @@ from __future__ import annotations +import uuid from contextlib import asynccontextmanager +from datetime import UTC, datetime from types import SimpleNamespace import pytest -from app.repository import Repository +from app.domain import safe_hash +from app.reconciliation import IncrementalReconciler +from app.repository import LeasedTask, Repository class FakeResult: @@ -16,9 +20,15 @@ class FakeResult: def __iter__(self): return iter(self.rows) + def mappings(self): + return self + def scalar_one_or_none(self): return self.scalar + def scalar_one(self): + return self.scalar + class FakeConnection: async def execute(self, statement): @@ -38,10 +48,29 @@ class FakeConnection: class FakeEngine: + def __init__(self) -> None: + self.queries: list[str] = [] + @asynccontextmanager async def connect(self): yield FakeConnection() + @asynccontextmanager + async def begin(self): + yield ClaimConnection(self.queries) + + +class ClaimConnection: + def __init__(self, queries: list[str]) -> None: + self.queries = queries + + async def execute(self, statement, params): + sql = str(statement) + self.queries.append(sql) + if "INSERT INTO bitrix_sync.workflow_instances" in sql: + return FakeResult(scalar=params["id"]) + return FakeResult() + @pytest.mark.asyncio async def test_status_uses_common_status_alias_for_workflows() -> None: @@ -55,3 +84,95 @@ async def test_status_uses_common_status_alias_for_workflows() -> None: assert result["commands"] == {"succeeded": 4} assert result["webhook_lag_seconds"] == 1.5 assert result["settings_version"] == 1 + + +@pytest.mark.asyncio +async def test_claims_recover_expired_leases() -> None: + repository = object.__new__(Repository) + engine = FakeEngine() + repository.engine = engine + + assert await repository.claim_tasks("worker", 10, 60) == [] + assert await repository.claim_webhooks("worker", 10, 60) == [] + + assert "status='leased' AND locked_until < now()" in engine.queries[0] + assert "status='processing' AND locked_until None: + repository = object.__new__(Repository) + engine = FakeEngine() + repository.engine = engine + task = LeasedTask(uuid.uuid4(), "contact.map_or_create", uuid.uuid4(), uuid.uuid4(), 0) + + workflow_id = await repository.create_workflow(task) + + assert isinstance(workflow_id, uuid.UUID) + assert "SET state='running'" in engine.queries[1] + + +@pytest.mark.asyncio +async def test_reconciliation_uses_unambiguous_external_id_alias() -> None: + queries: list[str] = [] + + class Connection: + async def execute(self, statement, params): + queries.append(str(statement)) + return FakeResult() + + class ReconciliationRepository: + @asynccontextmanager + async def transaction(self): + yield Connection() + + reconciler = IncrementalReconciler(ReconciliationRepository(), None, "ufCrm_1") + await reconciler._enqueue_changed( + [ + ("29406", "2026-08-20T13:21:00+00:00"), + ("29502", "2026-08-20T13:22:00+00:00"), + ] + ) + + assert "AS candidate(external_id,event_id)" in queries[0] + assert "w.external_entity_id=candidate.external_id" in queries[0] + assert "ON CONFLICT (receiver_type,event_id)" in queries[0] + + +@pytest.mark.asyncio +async def test_apply_crm_profile_hashes_without_database_digest() -> None: + calls: list[tuple[str, dict | None]] = [] + + class Connection: + async def execute(self, statement, params=None): + calls.append((str(statement), params)) + return FakeResult() + + class Engine: + @asynccontextmanager + async def begin(self): + yield Connection() + + repository = object.__new__(Repository) + repository.engine = Engine() + source_updated_at = datetime.now(UTC) + + await repository.apply_crm_profile( + uuid.uuid4(), + "29406", + full_name="Тестовый пользователь", + citizenship=None, + email=None, + source_updated_at=source_updated_at, + source="webhook", + ) + + profile_sql, profile_params = calls[1] + snapshot_sql, snapshot_params = calls[2] + assert "source_updated_at=:source_updated_at" in profile_sql + assert profile_params["source_updated_at"] == source_updated_at + assert "digest(" not in snapshot_sql + assert "CAST(:external_id AS varchar(128))" in snapshot_sql + assert "CAST(:source AS varchar(24))" in snapshot_sql + assert snapshot_params["full_name_hash"] == safe_hash("Тестовый пользователь") + assert snapshot_params["email_hash"] == safe_hash("") diff --git a/VM2_services/codebase/services/docker-compose.yml b/VM2_services/codebase/services/docker-compose.yml index fd5301c..cdf19f9 100644 --- a/VM2_services/codebase/services/docker-compose.yml +++ b/VM2_services/codebase/services/docker-compose.yml @@ -44,6 +44,7 @@ x-bitrix-sync-environment: &bitrix-sync-environment BITRIX_SYNC_CONTACT_USER_ID_FIELD: ${BITRIX_SYNC_CONTACT_USER_ID_FIELD:?set contact field} BITRIX_SYNC_CONTACT_REGISTERED_FIELD: ${BITRIX_SYNC_CONTACT_REGISTERED_FIELD:?set registration field} BITRIX_SYNC_CONTACT_CITIZENSHIP_FIELD: ${BITRIX_SYNC_CONTACT_CITIZENSHIP_FIELD:?set citizenship field} + BITRIX_SYNC_CONTACT_SOURCE: ${BITRIX_SYNC_CONTACT_SOURCE:?set contact source} BITRIX_SYNC_WEBHOOK_ALLOWED_CIDRS: ${BITRIX_WEBHOOK_ALLOWED_CIDRS:-} BITRIX_SYNC_HTTP_TIMEOUT_SEC: ${BITRIX_SYNC_HTTP_TIMEOUT_SEC:-10} BITRIX_SYNC_DB_POOL_SIZE: ${BITRIX_SYNC_DB_POOL_SIZE:-5} @@ -372,6 +373,7 @@ services: image: ${BITRIX_SYNC_IMAGE:?set immutable bitrix-sync image digest} command: ["han-bitrix-sync-reconciliation"] user: "10001:10001" + restart: "no" environment: *bitrix-sync-environment volumes: - *postgres-ca-volume diff --git a/VM2_services/documentation/module-07-bitrix-sync.md b/VM2_services/documentation/module-07-bitrix-sync.md index 0aba4d0..9bcfd0f 100644 --- a/VM2_services/documentation/module-07-bitrix-sync.md +++ b/VM2_services/documentation/module-07-bitrix-sync.md @@ -582,6 +582,7 @@ BITRIX_SYNC_MODE=full BITRIX_SYNC_CONTACT_USER_ID_FIELD=UF_CRM_... BITRIX_SYNC_CONTACT_REGISTERED_FIELD=UF_CRM_1778692456 BITRIX_SYNC_CONTACT_CITIZENSHIP_FIELD=UF_CRM_... +BITRIX_SYNC_CONTACT_SOURCE= BITRIX_SYNC_PORTAL_HOST= BITRIX_SYNC_PORTAL_MEMBER_ID= BITRIX_SYNC_PUBLIC_BASE_URL=https:// diff --git a/VM2_services/vm1-bug-v10.tar.gz b/VM2_services/vm1-bug-v10.tar.gz new file mode 100644 index 0000000000000000000000000000000000000000..d95f242d3fc45fbdbc5fa1966a3b512165a58746 GIT binary patch literal 3295 zcmV<53?TC#iwFRVABSoH1MORRZ`(K)pT7(AI|#v|Gh0NqeB`KhG7Cjk6QMpbvXf4C zPzbWb*g|DV)sfUY=y%^oQj&F0w%nxcbP*M`V;+x>kH2@3REUB*+eBg8^2`q?;7xDk zdm(DASI$q5^$3D+e0^i^-A4Z<59SvD-F1_tsvt(;+`PFlft1 zE2c~^$19kFu|Rxi+@fgBBlpe@Ja>g$p!lM`64R5eu5@L5gPBhMB+Mp12y-K4I+ZRI zeJZ_HX79+faC=rUFSN)%R==x>#_C6k`BUMP+3D>>d!uN&JQ>T^nx>59cMKAYM*IfgXw*n5 zU(5G;9j3q}ZKl4KrvH8|Y0`M6j+HXGY6C*^V{8P58`&=6<}r9HN_m1iDSu$`1K1|B;Z?>C zrn{}0Uq+EHw%cBWVAh`FAT#m;vES>QFkvdqq;qMi;92G+2-Yt^QJ6|yr&8A~BQoYj zi1avc7}PXnpiUc-3Xo-oHFTzmcB9IQKAXHz#tmqN)Z^=s_W`-(rP+ua_>%^@B+D?G z{ziGnq8#@0&6?=M1{x-IBKilk-bH$<$QmqibJDbmlE&tBZbyOri4WJVS+cD)HBAlc zxjNI-KlJIlu?%I>`A`|)_N-$5)^I{p>4g=u+=Z0!u#U#uK~|*;2B{MM9VHONS z;N5}s^SvN4oH|+~WpoahS1a6C>Iu~qTFp_w*$m#U&xfk4r`WHMD^(L}emr+4=q#z`ThP$xc zPf5!0>#rjBvXIqqoV@}siSaPLP{)7h>Ij@HRU~`0iX&{cDGl0^vuwy|Xbi$}F$^Vg z2OXg!b~=L4IiSUPzeGy_y4nOgRdsxkTF%ty*VNE7g<1DGC0Ys!uCJyu@Uo+XbSWM* zd^_!zmvyOA&WqC4*xQ8JGtH_}Ye&($%67lQ!+ zJqhV?qG&{jJ4#Ci!an&9O*Xtl|L(9ae&0U`;6E2V)~?G&Mv@4$dzzR#wR z)x~fe@7@8QQ0S8@&39_w8A02z=bNyc`b0i*pgYEJ?)UbNh6&j&DZanIZ=0Ytb1%yJ zB&v{;OPcaR^o;`7aC8_#$KC8(+l1MrEjJnN?v^--5IF(XcE>o$I;Jt;mT{OZ%i}~V z|299jHWY}#{h{ljPsoI?Ft#oDve4aO2tM@exW6jN(lpE^qKmbw=OGlGffgSw@Hf`v zL(I^baB?J+sKN?vrFPns*ak@$t1aunRAa|p**QSC-mUE*aUH1Q|XUL3qn2@$;B8IH(x zj>3$(y+U3L1N5>YP{6f0bPvOXf$6_N#m&^s$TI!@qg z`Z3V!v;^$Q@CcAQB$937#b|!X^kNz&wZ)6Qdi6?#|B2Z(z17NaEiBH22apH~i16YN zER|;(4iO`%rk@2dvJ7{*t+<0$!9wO`!Q)U$Y~#;e=8&e*NgiyNPSrf+3AYL@%4R92 z9!0f4pYR4cmez6aUoVOmvyU(y;J02Vj(<+SY4T$bZ4vXQrR z>~&LjAW(Nubp{Q9iZ?AA0+pX$D#Ia`E)jER8ZO84K45_BaKE*Iv9=pK4%hwdMTZ5| zo=}DAQ5C93Y%6-qvZ6hf6J+)WXD3oExLE zjkc!*w`TsC#H-2-I!d;u&PZs3nFJnj8JN3yvUSy38!IPW5Z1myldWs~Xv$E*yHNPB zOr4A=;)+G2U#T$bPb8u#BlxavlaCJ`!{P9$guQjDivrtFq5_(Wf zg0Ii@v7|T0A*nI99Xnc=w(5vj#a(hQxc5AMQ!APD&Eq*UDrwIG8zT2tCIRS+R^Zy! ziOtC+H;a3m1$zPM(WyE_|k|ASPdWG@)^-7D?d8wblK5qZSJ{M{Y(31O!Y5PX4!z8suvqqy?x0@uoRsz zy#0Ove#!^-bdOEfCU}oD+v2P0OQWM=v01CU+Pq|6#&}4?P*JfYg3T*et=6TIYoc-+ zwICas@9Q$j9Y?&(Muo#$W{b*^EYk)sHUMf209E64W1aTb(J~1x$)1S4!dtx@uneZJ z6(BG?Vf3%k?-u;@fXQEePL4=A-4k-lq(e4^+jLwoo_FZy!$k4Vy8nFqt(9)mt(B&x z149lq3Ld_)^+-7RZbpO-|Kr+`P0;?$49QnB89VP<)HciCn|NH)Qrr@1(j%{;{pNXl z88-GJfOzxUX=pL~G9a3p+pOd(_JJ`AVH*n5X5{)M#bu;fen2+sajJUY`Qf&NpnMN) zGJdUI=+a=Mj`zK*L_=Xn&sMg}gy>%9?JR3itJM|Nd{lMYF|oZ(N|tzj@rFV5kYg)= z^9Qs{^ycX8a4y&i=dVUVjJTI;;a+aJSoN?q-^Uy$Ogw5_x9w_1wME@!3r5S@%M;bM zb=M)PEgA9Rmjn4JKFl_iPW$ZmFc16Y!Qofb+%Jd=zoPO!il#p4WQ}v#E0f8uF>vj< zCM*iMzxbUITV&=-0;4A|{~v(4bF|y+3ayiBsI*Va|4gmsG`Rt{IiD<=oVr%kw5Te^^5G^aWSl-c zf0kdq>TMcaT0Jt(m8Pe|g#FEN3B7Z>S+Fg8AIUt)#WaeyY?CO(?v~eY7{qHcz5Pe3 z&7>K_W?MD;50kwC#2{O+!26RfT$?7y*tUJtfPS-=glb0G>_9zmq|G}+m`>fD)dqvy zjJmoyo+;W`8j^!YS-~&hsj@UvDx`DQRKJ9DdgD3ax-JrEK)UW_gECpRRors`wvQbr=uzX7zI`f# d#J~5e=ugkn^YlDDPtP}e{s$HSYXksB007nWbYcJi literal 0 HcmV?d00001 diff --git a/support¬es/backlog.md b/support¬es/backlog.md index ae75d5d..19d902a 100644 --- a/support¬es/backlog.md +++ b/support¬es/backlog.md @@ -77,6 +77,9 @@ 30. #INFRASTRUCTURE перевести взаимодействие с signoz на TLS (сейчас OTEL_REMOTE_TLS_INSECURE=true) 31. #BACK_BUSINESS VM1 -> VM2: `curl -sS --cacert /etc/han/ca/vm2-internal-ca.crt "https://processing.internal:8443/internal/safety/status" | python3 -m json.tool` В коде захардкожен Redis в components = "degraded". Надо реализовать реальную проверку вместо костыля. 32. #BACK_BUSINESS Решить проблему с обновлением сигнатур CLAMAV. +33. #BACK_BUSINESS Создать поле "Номер телефона в приложении". Подробности реализации ниже (на подумать) +34. #BACK_BUSINESS реализация смены номера телефона. Подробности реализации ниже (на подумать) +35. #BACK_BUSINESS схлапывание контактов. Подробности реализации ниже (на подумать) # Критично для релиза: ~~1. Разработка message-safety~~ @@ -85,3 +88,170 @@ ~~4. Подключить OTLP-провайдер~~ ~~5. Починить UI баги~~ 6. Второй контур для продакшн + + +33. details +Комплексный алгоритм: + +1. Нормализовать номер приложения `P`. +2. Если active mapping существует: + - Contact принадлежит текущему `user_id` — обновить служебные поля и app-телефон; + - `user_id` пуст — восстановить служебные поля; + - указан чужой `user_id` — ничего не перезаписывать, создать alert `mapping_identity_mismatch`; + - Contact отсутствует — пометить mapping broken и создать alert. +3. Если mapping отсутствует, собрать объединённый список Contact: + - найденные по текущему `user_id`; + - найденные по новому полю app-телефона `P`; + - найденные через `duplicate.findbycomm` по стандартному `PHONE`. +4. Загрузить все карточки и классифицировать: + - принадлежат текущему пользователю; + - свободны — `user_id` пуст; + - принадлежат другому пользователю; + - app-телефон равен `P`; + - app-телефон пуст — legacy-контакт; + - app-телефон отличается от `P`. + +Дальнейшие ветки: + +- Найден Contact с тем же `user_id`: + - выбрать самый новый среди таких Contact; + - восстановить mapping; + - записать `P` в app-телефон; + - если таких Contact несколько или есть другие совпадения — создать duplicate alert со всеми ID. + +- Контактов с тем же `user_id` нет: + - кандидатами считаются Contact с app-телефоном `P`; + - также допускаются legacy-контакты с пустым app-телефоном, найденные по стандартному `PHONE`; + - Contact с другим непустым app-телефоном нельзя автоматически привязывать. + +- Один кандидат свободен: + - записать `user_id`, registration flag и app-телефон; + - создать mapping. + +- Несколько кандидатов: + - выбрать самый новый по `CREATED_TIME`, затем по числовому ID; + - если выбранный свободен — привязать его и создать `duplicate_contacts`; + - если выбранный принадлежит другому пользователю — создать новый Contact и alert `contact_owned_by_other_user`; + - в `contactIds` alert передать все старые ID и ID нового Contact. + +- Совпадения есть только по стандартному `PHONE`, но у них указан другой app-телефон: + - не считать их Contact текущего пользователя; + - создать новый Contact; + - создать alert о совпадении общего номера. + +- Совпадений нет: + - создать новый Contact без alert. + +Дополнительные правила: + +- app-телефон — App-master: изменения этого поля в Б24 не меняют номер авторизации. +- При смене номера приложение обновляет app-телефон и гарантирует наличие нового номера в стандартном `PHONE`, но не удаляет остальные номера Contact. +- Перед повторным `contact.add` после timeout необходимо повторно искать по `user_id` и app-телефону, чтобы не создать дубль. +- Mapping создаётся только после успешной записи служебных полей. +- Все конфликтные ветки должны иметь отдельные регрессионные тесты до реализации. + + +34. details +Изменение стандартного `PHONE` в Битрикс24 не должно автоматически менять номер авторизации в приложении. Иначе сотрудник с доступом к CRM сможет случайно или намеренно передать чужой аккаунт другому номеру. + +Рекомендуемый процесс: + +1. Сотрудник меняет стандартный номер в карточке. +2. Webhook фиксирует расхождение с полем «Номер телефона в приложении». +3. Система создаёт отдельный запрос/смарт-процесс «Смена номера приложения». +4. Сотрудник явно выбирает, какой из номеров Contact является новым номером приложения. +5. Выполняются проверки: + - нормализация E.164; + - новый номер не занят другим пользователем; + - Contact действительно связан с нужным `user_id`; + - нет другого активного запроса. +6. Клиент подтверждает смену: + - OTP на новый номер; + - подтверждение в активной сессии приложения; + - желательно OTP на старый номер, если он доступен. +7. Если старый номер недоступен — отдельный усиленный сценарий идентификации сотрудником с обязательным аудитом и, желательно, дополнительным согласованием. +8. После подтверждения App меняет номер в своём источнике истины — БД приложения/Keycloak. +9. Через outbox создаётся задача синхронизации. +10. Bitrix Sync обновляет: + - поле «Номер телефона в приложении»; + - служебные поля связи; + - при необходимости добавляет новый номер в стандартный `PHONE`; + - закрывает запрос на смену номера. + +Важные правила: + +- Обычное редактирование `PHONE` только создаёт предложение на смену, но не меняет логин. +- Поле «Номер телефона в приложении» остаётся read-only для сотрудников. +- Старые дополнительные номера Contact автоматически не удаляются. +- Нужны срок действия запроса, ограничение попыток OTP, журнал сотрудника и времени операции. +- Повторный запрос должен быть идемпотентным. + +Так сохраняется удобство работы через Б24, но CRM не становится небезопасным источником истины для номера авторизации. + +35. details +Здесь нужно объединять не только карточки Битрикс24, но прежде всего два аккаунта приложения. Простое слияние Contact оставит данные личного кабинета привязанными к старому `user_id`. + +Рекомендуемый сценарий account merge. + +1. Создать заявку на объединение: + - старый `user_id` и старый Contact; + - новый `user_id` и новый Contact; + - результат комплайенса; + - оператор и подтверждающий сотрудник; + - причина и audit trail. + +2. Выбрать канонический аккаунт. Обычно это старый аккаунт, содержащий историю и данные. Новый пустой аккаунт становится поглощаемым. + +3. На время операции заблокировать изменения обоих аккаунтов и повторные merge-заявки. + +4. Проверить конфликты: + - заказы, документы, согласия; + - бонусы, баланс и другие финансовые данные; + - активные процессы; + - наличие данных в новом аккаунте; + - отсутствие третьего аккаунта с новым номером. + +5. На стороне приложения: + - привязать новый номер к каноническому старому аккаунту; + - старый номер сделать неактивным либо историческим; + - перенести допустимые данные нового аккаунта; + - отметить новый `user_id` как `merged_into=`; + - запретить дальнейшее использование поглощённого аккаунта; + - обновить Keycloak/механизм авторизации; + - отозвать старые сессии и потребовать повторный вход. + +6. Только после успешного App merge обработать Битрикс: + - старый Contact выбрать каноническим; + - записать в него канонический `user_id`; + - поле «Номер телефона в приложении» установить в новый номер; + - объединить стандартные номера обоих Contact; + - создать active mapping канонического пользователя на канонический Contact; + - у второго Contact очистить app-телефон, `user_id` и registration flag; + - затем объединить его с каноническим Contact штатным механизмом Б24 либо архивировать как поглощённый. + +7. Закрыть старые mapping со специальной причиной `account_merge`, но сохранить историю. + +8. Завершить заявку только после проверок: + - вход по новому номеру открывает старый личный кабинет; + - существует один active App user; + - существует один active mapping; + - только один Contact содержит app-телефон и registration flag; + - поглощённый аккаунт больше не создаёт задачи синхронизации. + +Ключевой принцип: сначала объединение аккаунтов приложения, затем CRM. Если сначала слить карточки Б24, это не вернёт пользователю данные. + +Для устойчивой реализации потребуется отдельная audited-процедура наподобие: + +```text +request_account_merge( + canonical_user_id, + absorbed_user_id, + canonical_contact_id, + absorbed_contact_id, + reason, + operator_id, + compliance_case_id +) +``` + +Операция должна быть идемпотентной, поэтапной и восстанавливаемой после сбоя. Для финансовых или юридически значимых данных автоматический перенос без отдельных правил недопустим. \ No newline at end of file