Set CRM order status to PACKED after Checkbox accepts the receipt
- ExoCrmClient.set_status: live SetStatus needs {Orders: [id], Status} and
replies per order, unlike the documented {ID, Status}
- receipts.crm_status_set_at (migration 0006); set right after creation,
retried by cron, row-locked to avoid a repeat PACKED overwriting a newer status
- CRM errors under capitalized 'Errors' and non-JSON replies are reported
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
@@ -36,6 +36,7 @@ class AuditAction(str):
|
||||
CASH_REGISTER_UPDATED = "cash_register.updated"
|
||||
RECEIPT_CREATE_REQUESTED = "receipt.create_requested"
|
||||
RECEIPT_CANCELLED = "receipt.cancelled"
|
||||
ORDER_CRM_STATUS_SET = "order.crm_status_set"
|
||||
|
||||
|
||||
class AuditLog(UUIDPrimaryKeyMixin, Base):
|
||||
|
||||
@@ -70,6 +70,9 @@ class Receipt(UUIDPrimaryKeyMixin, TimestampMixin, Base):
|
||||
# Снимок отправленного в Checkbox тела — для разбора спорных случаев.
|
||||
request_body: Mapped[dict[str, Any]] = mapped_column(JSONB, nullable=False, default=dict)
|
||||
last_checked_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True))
|
||||
# Когда заказу в CRM выставлен статус после создания чека (PACKED). NULL при
|
||||
# живом чеке — ещё не выставлен (CRM была недоступна), cron повторит.
|
||||
crm_status_set_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True))
|
||||
|
||||
__table_args__ = (
|
||||
# Не больше одного «живого» чека на заказ — защита от двойного нажатия
|
||||
|
||||
@@ -13,3 +13,5 @@ class CrmError(Exception):
|
||||
|
||||
class CrmClient(Protocol):
|
||||
async def get_orders(self, *, status: str) -> list[OrderOut]: ...
|
||||
|
||||
async def set_status(self, *, order_id: str, status: str) -> None: ...
|
||||
|
||||
@@ -1,7 +1,9 @@
|
||||
"""Реальный клиент exoCRM (`GetOrders`)."""
|
||||
"""Реальный клиент exoCRM (`GetOrders`, `SetStatus`)."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import Any
|
||||
|
||||
import httpx
|
||||
|
||||
from app.core.config import Settings
|
||||
@@ -10,6 +12,14 @@ from app.services.crm.checksum import compute_md5sum
|
||||
from app.services.crm.client import CrmError
|
||||
|
||||
|
||||
def _errors_text(data: dict[str, Any]) -> str:
|
||||
# Ключ бывает и `errors` (dict код → текст), и `Errors` (список строк).
|
||||
errors = data.get("errors") or data.get("Errors") or {}
|
||||
if isinstance(errors, dict):
|
||||
return "; ".join(f"{code}: {text}" for code, text in errors.items())
|
||||
return "; ".join(str(error) for error in errors)
|
||||
|
||||
|
||||
class ExoCrmClient:
|
||||
def __init__(self, settings: Settings) -> None:
|
||||
self._base_url = settings.crm_base_url
|
||||
@@ -18,29 +28,51 @@ class ExoCrmClient:
|
||||
self._shop_key = settings.crm_shop_key
|
||||
self._sid = settings.crm_sid
|
||||
|
||||
async def _post(self, method: str, params: dict[str, Any]) -> dict[str, Any]:
|
||||
body = {"apikey": self._api_key, "object": "Orders", "method": method, "params": params}
|
||||
body["md5sum"] = compute_md5sum(body, self._secret_key)
|
||||
|
||||
async with httpx.AsyncClient(timeout=30) as client:
|
||||
response = await client.post(self._base_url, json=body)
|
||||
response.raise_for_status()
|
||||
try:
|
||||
data = response.json()
|
||||
except ValueError as exc:
|
||||
# PHP-notice'ы CRM перед JSON — признак неверных параметров.
|
||||
raise CrmError(f"CRM вернула не JSON: {response.text[:300]}") from exc
|
||||
|
||||
if not isinstance(data, dict):
|
||||
raise CrmError(f"Неожиданный ответ CRM: {str(data)[:300]}")
|
||||
return data
|
||||
|
||||
async def get_orders(self, *, status: str) -> list[OrderOut]:
|
||||
body = {
|
||||
"apikey": self._api_key,
|
||||
"object": "Orders",
|
||||
"method": "GetOrders",
|
||||
"params": {
|
||||
data = await self._post(
|
||||
"GetOrders",
|
||||
{
|
||||
"sid": self._sid,
|
||||
"key": self._shop_key,
|
||||
"Status": status,
|
||||
"ReturnGoods": True,
|
||||
"ReturnTotals": True,
|
||||
},
|
||||
}
|
||||
body["md5sum"] = compute_md5sum(body, self._secret_key)
|
||||
|
||||
async with httpx.AsyncClient(timeout=30) as client:
|
||||
response = await client.post(self._base_url, json=body)
|
||||
response.raise_for_status()
|
||||
data = response.json()
|
||||
|
||||
)
|
||||
if data.get("status") != "OK":
|
||||
errors = data.get("errors") or {}
|
||||
message = "; ".join(f"{code}: {text}" for code, text in errors.items())
|
||||
raise CrmError(f"CRM вернула ошибку: {message or 'неизвестная ошибка'}")
|
||||
raise CrmError(f"CRM вернула ошибку: {_errors_text(data) or 'неизвестная ошибка'}")
|
||||
return [OrderOut.model_validate(order) for order in data.get("result") or []]
|
||||
|
||||
return [OrderOut.model_validate(order) for order in data.get("result", [])]
|
||||
async def set_status(self, *, order_id: str, status: str) -> None:
|
||||
"""`SetStatus`. Формат проверен на боевой CRM и расходится с документацией.
|
||||
|
||||
Документация показывает `params: {"ID": ..., "Status": ...}`, но CRM на это
|
||||
отвечает «Undefined order list.». Рабочий вариант — список ID в `Orders`,
|
||||
а ответ — результат по каждому заказу без общего `status: OK`:
|
||||
{"123901": {"Status": "Success", "ChangeStatus": "Success"}}
|
||||
"""
|
||||
data = await self._post("SetStatus", {"Orders": [order_id], "Status": status})
|
||||
result = data.get(order_id)
|
||||
if isinstance(result, dict) and result.get("Status") == "Success":
|
||||
return
|
||||
details = _errors_text(data) or (
|
||||
str(result) if result is not None else "нет ответа по заказу"
|
||||
)
|
||||
raise CrmError(f"CRM не сменила статус заказа {order_id} на {status}: {details}")
|
||||
|
||||
@@ -40,6 +40,11 @@ _FIXTURE_ORDERS: list[dict] = [
|
||||
class StubCrmClient:
|
||||
def __init__(self, orders: list[dict] | None = None) -> None:
|
||||
self._orders = orders if orders is not None else _FIXTURE_ORDERS
|
||||
# order_id → статус, выставленный через set_status (для проверок в тестах).
|
||||
self.statuses: dict[str, str] = {}
|
||||
|
||||
async def set_status(self, *, order_id: str, status: str) -> None:
|
||||
self.statuses[order_id] = status
|
||||
|
||||
async def get_orders(self, *, status: str) -> list[OrderOut]:
|
||||
return [
|
||||
|
||||
@@ -20,6 +20,7 @@ from datetime import UTC, datetime, timedelta
|
||||
from decimal import ROUND_HALF_UP, Decimal
|
||||
from typing import Any
|
||||
|
||||
import httpx
|
||||
from fastapi import Request
|
||||
from sqlalchemy import select
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
@@ -39,6 +40,7 @@ from app.services.checkbox.client import (
|
||||
CheckboxError,
|
||||
CheckboxUnavailableError,
|
||||
)
|
||||
from app.services.crm.client import CrmClient, CrmError
|
||||
|
||||
log = get_logger(__name__)
|
||||
|
||||
@@ -65,6 +67,16 @@ _CHECKBOX_TO_STATUS = {
|
||||
# или Checkbox был недоступен) и повторить его из cron'а.
|
||||
_PENDING_RETRY_AFTER = timedelta(minutes=1)
|
||||
|
||||
# Статус заказа в CRM, когда Checkbox принял ЕТТН-чек: заказ можно собирать.
|
||||
CRM_STATUS_AFTER_RECEIPT = "PACKED"
|
||||
# Чеки, при которых заказу нужен этот статус в CRM (Checkbox чек принял).
|
||||
_CRM_STATUS_RECEIPT_STATUSES = (
|
||||
ReceiptStatus.CREATED,
|
||||
ReceiptStatus.DONE,
|
||||
ReceiptStatus.RECEIPT_ERROR,
|
||||
ReceiptStatus.RETURNED,
|
||||
)
|
||||
|
||||
|
||||
class ReceiptValidationError(Exception):
|
||||
"""Заказ нельзя отправить в Checkbox — сообщение показывается кассиру."""
|
||||
@@ -500,6 +512,54 @@ async def sync_ettn_statuses(session: AsyncSession, client: CheckboxClient) -> N
|
||||
await session.commit()
|
||||
|
||||
|
||||
async def sync_crm_statuses(session: AsyncSession, crm: CrmClient) -> None:
|
||||
"""Переводит в CRM заказы с принятым Checkbox чеком в статус PACKED.
|
||||
|
||||
Отдельно от создания чека и идемпотентно по `crm_status_set_at`: сбой CRM
|
||||
не должен ни откатывать уже созданный в Checkbox чек, ни теряться —
|
||||
вызывается сразу после создания и повторяется cron'ом до успеха.
|
||||
"""
|
||||
receipt_ids = list(
|
||||
await session.scalars(
|
||||
select(Receipt.id)
|
||||
.where(Receipt.status.in_(_CRM_STATUS_RECEIPT_STATUSES))
|
||||
.where(Receipt.crm_status_set_at.is_(None))
|
||||
)
|
||||
)
|
||||
for receipt_id in receipt_ids:
|
||||
# Строка блокируется на время вызова CRM: задача создания и cron могут
|
||||
# сработать одновременно, а повторный PACKED откатил бы статус, который
|
||||
# менеджер уже успел сменить дальше. Занято или уже выставлено — пропуск.
|
||||
receipt = await session.scalar(
|
||||
select(Receipt)
|
||||
.where(Receipt.id == receipt_id)
|
||||
.where(Receipt.crm_status_set_at.is_(None))
|
||||
.with_for_update(skip_locked=True)
|
||||
)
|
||||
if receipt is None:
|
||||
continue
|
||||
order_id = receipt.order_id # после rollback атрибуты истекают
|
||||
try:
|
||||
await crm.set_status(order_id=order_id, status=CRM_STATUS_AFTER_RECEIPT)
|
||||
except (CrmError, httpx.HTTPError) as exc:
|
||||
await session.rollback() # снять блокировку строки
|
||||
log.warning("crm_status_set_failed", order_id=order_id, error=repr(exc))
|
||||
continue
|
||||
receipt.crm_status_set_at = datetime.now(UTC)
|
||||
await audit.record(
|
||||
session,
|
||||
action=AuditAction.ORDER_CRM_STATUS_SET,
|
||||
actor_label="worker",
|
||||
entity_type="order",
|
||||
entity_id=receipt.order_id,
|
||||
payload={"status": CRM_STATUS_AFTER_RECEIPT, "receipt_id": str(receipt.id)},
|
||||
)
|
||||
# Коммит на каждый заказ: статус в CRM уже сменён, отметку нельзя терять
|
||||
# из-за сбоя на следующем заказе.
|
||||
await session.commit()
|
||||
log.info("crm_status_set", order_id=receipt.order_id, status=CRM_STATUS_AFTER_RECEIPT)
|
||||
|
||||
|
||||
async def latest_receipts_by_order(
|
||||
session: AsyncSession, order_ids: list[str]
|
||||
) -> dict[str, Receipt]:
|
||||
|
||||
@@ -3,7 +3,8 @@
|
||||
Запускается отдельным процессом: `arq app.worker.WorkerSettings`.
|
||||
- `create_ettn_receipt` — задача, которую ставит API после запроса кассира;
|
||||
- `poll_np_statuses` — раз в минуту статусы ТТН по заказам без чека;
|
||||
- `poll_receipts` — раз в минуту повтор зависших `pending` и статусы `created`-чеков.
|
||||
- `poll_receipts` — раз в минуту повтор зависших `pending`, статусы `created`-чеков
|
||||
и повтор смены статуса заказа в CRM (PACKED), если CRM была недоступна.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
@@ -19,6 +20,7 @@ from app.core.logging import configure_logging, get_logger
|
||||
from app.db.session import SessionFactory
|
||||
from app.services import receipts as receipts_service
|
||||
from app.services.checkbox.client import get_checkbox_client
|
||||
from app.services.crm.exo_client import ExoCrmClient
|
||||
from app.services.nova_poshta.np_client import NpTrackingClient
|
||||
from app.services.orders import sync_np_statuses
|
||||
|
||||
@@ -29,6 +31,7 @@ async def startup(ctx: dict[str, Any]) -> None:
|
||||
configure_logging()
|
||||
ctx["np_client"] = NpTrackingClient(settings)
|
||||
ctx["checkbox_client"] = get_checkbox_client()
|
||||
ctx["crm_client"] = ExoCrmClient(settings)
|
||||
log.info("worker_starting", environment=settings.environment)
|
||||
|
||||
|
||||
@@ -43,12 +46,15 @@ async def create_ettn_receipt(ctx: dict[str, Any], receipt_id: str) -> None:
|
||||
await receipts_service.create_ettn_for_receipt(
|
||||
session, ctx["checkbox_client"], uuid.UUID(receipt_id)
|
||||
)
|
||||
# Чек принят Checkbox — сразу переводим заказ в CRM в PACKED.
|
||||
await receipts_service.sync_crm_statuses(session, ctx["crm_client"])
|
||||
|
||||
|
||||
async def poll_receipts(ctx: dict[str, Any]) -> None:
|
||||
async with SessionFactory() as session:
|
||||
await receipts_service.retry_pending_receipts(session, ctx["checkbox_client"])
|
||||
await receipts_service.sync_ettn_statuses(session, ctx["checkbox_client"])
|
||||
await receipts_service.sync_crm_statuses(session, ctx["crm_client"])
|
||||
log.info("receipts_polled")
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user