20 KiB
ARCHITECTURE (AS BUILT) — «Контракт-чек»
This document describes the system as it exists in the code today, verified against
src/,migrations/,tests/, and compose files. The design-time rationale, decision history and future steps remain indocs/ARCHITECTURE.md(design doc); this file is the as-built map. Domain language is defined in the rootCONTEXT.md.
1. What the system is
LLM-powered screening of contract (deal) documents under Russian / Belarusian civil law. A user (Telegram bot, web JWT auth, or B2B API key) uploads a document; an event-driven pipeline extracts text, prescreens it, optionally runs deep LLM analysis, and produces a markdown report with findings, quotes and clause references. Document consumption is metered: subscription Quota first, then prepaid Credits, with ЮKassa payments (topups, subscriptions, renewals, refunds with clawback).
Stack: Python 3.14, uv + hatchling, FastAPI, SQLAlchemy 2 (async) +
Alembic, aio-pika (RabbitMQ), MinIO (S3), Redis, aiogram 3, httpx,
Prometheus + structlog + Sentry, OpenObserve/Vector (or Grafana/Loki) for
observability, Nginx + certbot at the edge. Docker Compose deployment
(one Dockerfile per service in srv/), k8s-ready but not k8s-running.
2. Hexagonal shape and import rules
src/contract_check/
├── core/ # domain — owns ALL state and side effects
│ ├── db/ # models, session, repositories/
│ ├── mq/ # topology, publisher, consumer, messages, management
│ ├── s3/ # Storage port + MinIO adapter
│ ├── llm/ # LLMProvider port + ollama_cloud / yandex_gpt adapters
│ ├── billing/ # PaymentProvider port + yookassa, quota, fulfillment, refunds
│ ├── extraction/ # DocumentExtractor port + format adapters
│ ├── analysis/ # chunker, checklist, analyzer, report_schema
│ ├── notifications/ # SMTP transport + notify publisher
│ ├── security/ # argon2 passwords
│ └── config.py logging.py metrics.py telemetry.py sentry.py errors.py
│ credits.py tokens.py api_keys.py rate_limit.py auth*.py passkeys.py
│ redis_client.py review/
├── api/ # FastAPI image (routes/, admin/, billing/ pay pages, schemas/)
├── worker_extract/ # CPU image: pymupdf + mammoth + striprtf + chardet + tesseract
├── worker_prescreen/# hybrid metadata extraction + routing
├── worker_analyze/ # I/O image: LLM provider
├── worker_notify/ # email (aiosmtplib)
├── worker_billing/ # timer scheduler (no MQ)
└── bot/ # aiogram adapter — HTTP-only
Boundary rules (statically enforced for the bot by
tests/unit/test_bot_boundary.py):
core/*may import anything.api/*,worker_*/*import onlycore/*(+ their own modules).bot/*imports onlyhttpx/aiohttp/aiogram and its own modules — nevercore.db,core.s3,core.llm,core.mq,core.credits. It talks to the api over HTTP exclusively, and may run on a separate host (docker-compose.bot.yml).
Control plane: api is the only writer for user-initiated mutations
(upload → reserve Document Slot → MinIO → publish). Workers write their own
stage rows / reports; credits move only through the idempotent compensation
helpers (core/billing/quota.compensate_document_slot, core/credits.refund_credit).
3. Services and the pipeline
Telegram ──► bot (aiogram, HTTP-only) ──HTTP──► api (FastAPI) ──publish──► RabbitMQ
owns: PG, MinIO, Redis, │
quota/credits, tokens ▼
extract.q ──► worker-extract
(CPU: pymupdf/mammoth/tesseract)
│ publish
▼
prescreen.q ──► worker-prescreen
(hybrid heuristic + optional LLM)
│ deep_analysis │ manual_review (terminal)
▼ ▼
analyze.q ──► worker-analyze ──► Postgres (Report, done)
notify.x ──► notify.q ──► worker-notify (SMTP: reset, magic link, …)
worker-billing (60 s Postgres-advisory-lock tick; no MQ)
api (api/, srv/api/)
create_app(api/app.py:115): logging/Sentry, MinIOensure_bucket, MQ publisher (publisher confirms), Redis rate-limiter with in-memory fallback, admin seed. Middleware: CORS, correlation-id propagation,http_request_duration, access logs, redacted debug payloads, generic 500.- Routes under
/api/v1:auth/(Telegram bot/web/miniapp, email+password JWT pair, passkeys, magic links, support),documents(upload → 202),reports(polling + SSE/events),me(profile, overview, documents, telegram bind, password),billing(plans, checkout, invoices, subscription, autorenew, refund),b2b(API-key analyze/status/usage/ profile/keys CRUD),webhooks/yookassa. Plus/healthz,/readyz,/metrics(bearer-gated), server-rendered/pay/{invoice_id}(short-lived JWT pay token), and/admin/*(Jinja2 + HTMX panel: users, credits, invoices/refunds, hold clearing, gift subscriptions, manual-review queue, MQ queues/DLQ requeue). - Upload flow (
api/services.py:56upload_and_enqueue): validate suffix/format → reject onbilling_hold(402) → stream body with size cap (413) → MinIO put (DB failure triggers best-effort S3 cleanup) → insertdocuments(queued) +jobs(extract/pending) → reserve Document Slot (Quota-then-Credit whenPLANS_ENABLED, else legacyreserve_credit; 402 onNoCredits) → publishDocumentUploaded(publisher confirms; on MQ failure: release slot / refund credit, statuspublish_failed) →202 {document_id, correlation_id, credits_left}.
worker-extract (worker_extract/, metrics :9101, prefetch 1)
Consumes extract.q / DocumentUploaded. Idempotency check against terminal
statuses → status extracting → download blob → extract_document
(core/extraction/factory.py: python-magic sniff → uploader MIME → suffix;
adapters: pdf/pymupdf, docx/mammoth, rtf/striprtf, txt+csv/chardet,
images/tesseract OCR; unsupported formats are terminal + refunded) → upload
extracted markdown to MinIO → publish PrescreenRequested (or
DocumentExtracted straight to analyze.q when PRESCREEN_ENABLED=false).
Failure classes: extraction_failed (non-refundable under infra_only),
ocr_failed, infra. Terminal failure → DLQ + slot compensation.
worker-prescreen (worker_prescreen/, metrics :9104, prefetch 1)
Consumes prescreen.q / PrescreenRequested. Handler
(worker_prescreen/handler.py:116): transition prescreening → download
text → hybrid extractor (extractor_hybrid.py): Stage 1 deterministic
heuristic (heuristic-v2, regex-free, keyword/positional); Stage 2 LLM
fallback only if enabled AND confidence < PRESCREEN_LLM_FALLBACK_THRESHOLD
(LLM failure never fatal; LLM fills heuristic Nones, booleans OR-merged,
confidence rescored) → router (router.py:30): missing
type/parties/low-confidence → manual_review (terminal status, admin queue);
amount ≥ PRESCREEN_HIGH_VALUE_THRESHOLD or penalty/arbitration clause →
deep_analysis (publish AnalyzeRequested with prescreen_meta);
auto_approve disabled by default (remapped to manual_review; when enabled
writes a lightweight Report, status=done). Persists prescreen_results.
worker-analyze (worker_analyze/, metrics :9102, prefetch 3)
Consumes analyze.q / AnalyzeRequested. Status analyzing/llm → download
text → provider.analyze(text, checklist) (chunking, fan-out under
asyncio.Semaphore(LLM_MAX_CONCURRENCY), json-schema + repair loop, 429 →
fallback model — all inside the LLM adapter) → render_markdown (findings,
quotes, clause refs, «не заменяет юриста» disclaimer) → reports upsert,
status=done. Failure classes: llm_quota, llm_timeout,
llm_invalid_output, infra; DLQ → slot compensation.
worker-notify (worker_notify/, metrics :9103, prefetch 5)
Consumes notify.q / NotificationMessage (kinds: password_reset, welcome,
email_verification, magic_link). SMTP via aiosmtplib; dev logger when
SMTP_HOST empty.
worker-billing (worker_billing/, metrics :9105, no MQ)
60-second tick guarded by a Postgres advisory lock (single replica):
renewal invoices for auto-renew subscriptions expiring < 3 days; period roll /
past_due (grace) / expired transitions; reconciliation of pending
invoices older than 15 minutes via the same apply_payment_status used by
the webhook. No-op safe when ЮKassa disabled.
bot (bot/, srv/bot/)
aiogram 3, polling or webhook mode (aiohttp, secret-token auth). Commands:
/start /help /profile /balance /plans /reports /status; checkout callbacks
(plan:*, topup:*); document upload → API → poll report with stage labels;
per-user rate limit (Redis, memory fallback); service-token login cached per
telegram_id; ensure_disclaimer guarantees the legal disclaimer.
4. RabbitMQ topology (declared idempotently in core/mq/topology.py)
| Object | Type / args | Purpose |
|---|---|---|
contracts.x |
direct | main exchange; RK extract / prescreen / analyze |
extract.q / prescreen.q / analyze.q |
quorum, DLX → contracts.retry.x |
main work queues |
contracts.retry.x + *.retry.q |
direct + classic (per-message TTL, lazy policy via mgmt API) |
delay slots; TTL expiry dead-letters the message back to contracts.x |
extract.dlq / prescreen.dlq / analyze.dlq |
quorum | poison; manual requeue via admin panel / core/mq/management.py |
notify.x / notify.q / notify.retry.q / notify.dlq |
mirror of the above | notification pipeline |
Retry mechanics (core/mq/consumer.py): generic Consumer[MsgT] with
handle/classify/on_failure/on_dlq hooks. On failure the consumer
re-publishes to the retry exchange with expiration = MQ_RETRY_BASE_MS · 2^(attempt−1) (default 2 s base → 2/4/8/16/32 s) and acks the original.
x-attempt >= MQ_MAX_ATTEMPTS (5) → DLQ with x-failure-class headers,
jobs.dlq=true, documents.status=failed, slot compensation. Pydantic
validation failure → straight to DLQ (infra poison). core.errors.TerminalError
bypasses retries. If the failure hooks themselves raise, the message is
retried (hook-failure fallback). Publisher uses publisher confirms —
unconfirmed publish fails the request, so a paid slot never silently vanishes.
Message schemas (core/mq/messages.py, pydantic v2, base PipelineMessage
with next_attempt()): DocumentUploaded, DocumentExtracted,
PrescreenRequested, AnalyzeRequested (carries prescreen_meta),
PrescreenCompleted, NotificationMessage. Headers: x-correlation-id,
x-attempt, x-origin; delivery_mode=2.
Idempotency: every handler re-reads documents.status first; terminal or
in-flight → ack and exit. No double-LLM, no double-refund.
5. Data model (Postgres, 16 tables, core/db/models.py; migrations 0001–0011)
| Table | Role |
|---|---|
users |
identity (telegram_id and/or email), argon2 password_hash, reset/magic-link token hashes, role (user/admin), credits_left (CHECK ≥ 0), billing_hold |
passkey_credentials |
WebAuthn credentials (credential_id, public_key, sign_count) |
documents |
one upload: s3 keys, status queued→extracting→prescreening→ocr→analyzing→done|failed|manual_review, refunded |
reports |
1:1 with document: content_json JSONB + markdown, model, tokens, latency, prescreen link |
prescreen_results |
extracted metadata (type, parties, amount, dates, clause flags), confidence_score, routing_decision, extractor_version |
jobs |
per-stage correlation (UNIQUE document+queue), attempts, last_failure_class, dlq |
service_tokens |
per-adapter bearer auth (bot/web/cli), revocable |
api_keys / api_key_requests |
B2B keys (SHA-256 hash, per-key rps + monthly quota) and per-call usage ledger |
invoices |
money in integer kopecks; kinds topup/subscription/renewal; statuses draft→pending→succeeded/cancelled/refunded |
plans / subscriptions / quota_usage |
seeded plan catalog; one active-or-past_due subscription per user (partial unique); quota ledger with UNIQUE(document_id) idempotency |
credit_events |
append-only credit ledger (delta ≠ 0, balance_after ≥ 0 self-verifying) |
user_profiles |
1:1 passive settings (language, TZ, notif/dashboard prefs) |
Data access: all SQL lives in core/db/repositories/ (14 repositories);
a repository receives an AsyncSession and never commits — callers own
transactions. Raw text() only inside repositories/migrations. Status/enum
columns are TEXT + CHECK (not PG enums) so migrations stay additive. Alembic
is async; CI exercises upgrade head → downgrade base.
6. Credits, Quota, refunds (billing invariants, core/billing/)
- Reserve-on-enqueue: the API synchronously reserves a Document Slot
before publishing (
core/billing/quota.py:22reserve_document_slot):SELECT … FOR UPDATEon the active subscription → quota rows this period <plan.monthly_quota→ insertquota_usage(ON CONFLICT DO NOTHING) → sourcequota; else atomic credit reserve →credits; else 402.PLANS_ENABLED=false→ legacy credits-only path. - Compensation (exactly-once): on terminal processing failure,
compensate_document_slot(quota.py:100) releases the quota row or refund the credit (guarded bydocuments.refunded), honouringREFUND_POLICY=all|infra_onlywith the failure-class taxonomy (extraction_failednon-refundable underinfra_only). - Payments:
PaymentProviderport (port.py) +YookassaProvider(httpx, Basic auth, Idempotence-Key = invoice UUID, kopecks↔RUB only at the HTTP boundary).YOOKASSA_ENABLED=false→PaymentsDisabled→ billing mutations 503, catalog readable. - Webhook trust (
api/routes/webhooks.py): Basic auth constant-time check, then the payment is re-fetched via REST — the push payload's amount is never trusted; fulfillment is idempotent (core/billing/fulfillment.apply_payment_status), also used by the reconciliation sweep. - Refunds (D10) (
core/billing/refunds.py): withinREFUND_WINDOW_DAYS(14) and usage ≤ 20 % of purchased → full; else proportionalmax(0, amount − used × PRICE_PER_DOC); otherwise 409. Execution: provider refund → invoicerefunded→ clawback (unspent credits / subscription + period quota rows) →users.billing_hold = TRUEwhen the balance would go negative; uploads 402 until an admin clears the hold.
7. LLM provider port (core/llm/)
LLMProvider Protocol (port.py:31): analyze(text, *, checklist) and
extract_prescreen(text). Registered providers via LLM_PROVIDER:
ollama_cloud (default; format: json-schema chat, model
qwen2.5:14b, fallback qwen2.5:7b) and yandex_gpt
(yandexgpt-lite, Api-Key, JSON_OBJECT). Shared behaviour in each adapter:
chunk fan-out under a semaphore, repair loop (one REPAIR re-send on invalid
JSON, then LLMError), 429 → fallback model → LLMQuotaError, exponential
HTTP backoff for timeouts/5xx, terminal LLMConfigError for connect/404.
Checklist: 10 frozen items (core/analysis/checklist.py); Finding schema
(report_schema.py) doubles as the Ollama json-schema and tolerates alias
field names.
8. AuthN/AuthZ surface
| Caller | Mechanism |
|---|---|
| bot adapter | Authorization: Bearer <service_token> (only /api/v1/auth/telegram/bot) → user JWT exchange |
| Telegram web / Mini App | Login Widget hash / initData verification → JWT |
| web user | email+password → JWT pair (access 24 h + refresh 30 d, Redis-backed revocation); passkeys (WebAuthn, Redis challenges) and magic links → single access JWT |
| B2B | X-API-Key (SHA-256, token-bucket rate limit, monthly quota) |
| admin panel | HttpOnly cookie + users.role = admin |
| ЮKassa webhook | Basic (shopId:secret) |
| health/metrics | none / bearer |
Redis is used only for rate limiting, refresh-token store, passkey challenges and bot limiter — never as a job queue.
9. Storage (MinIO, core/s3/)
Single bucket contract-check-docs; keys users/{uid}/docs/{did}.{ext}
(original) and users/{uid}/docs/{did}.txt (extracted markdown), built only
via key builders. Uploads proxy through the API (multipart, 25 MiB cap,
streamed with limit). minio-init one-shot creates the bucket and the ILM
expiry (DOC_RETENTION_DAYS 7 / TEXT_RETENTION_DAYS 30) — 152-ФЗ lever:
raw text leaves on schedule; reports live in Postgres and survive the purge.
10. Observability & deploy
- Logs/metrics/traces: structlog JSON + correlation-id contextvar
(propagated through HTTP headers and MQ headers), Prometheus per-service
(
/metricson api :9100; workers :9101–:9105 scraped in-network by Vector), Sentry (DSN-gated). Passive collection:observerprofile = OpenObserve + Vector;obsprofile = Prometheus + Loki + Promtail + Grafana. OTLP push was removed. - Compose (
docker-compose.yml): default = infra only (postgres:18 withwal_level=replica+ WAL archive, redis:8, rabbitmq:4, minio + minio-init); profilesservices(api + 5 workers),bot(standalone-able),edge(nginx + certbot + exporter),observer/obs. Host ports offset: PG 15432, Redis 17379, Rabbit 5672/15672, MinIO 9000/9001, api 8000, OpenObserve 5080, edge 80/443. - Images: one Dockerfile per service in
srv/(two-stage uv build,uv sync --no-default-groups --group <grp>→ python:3.14-slim); dependency groups (PEP 735) per service keep images minimal; bot is leanest. - CI (
.github/workflows/ci.yml): ruff check+format,tytypecheck, unit tests with--cov-fail-under=50, integration tests against an isolated compose test stack (docker-compose.test.yml, offset ports).
11. Testing
- Unit (
tests/unit/, ~40 files / 326 tests): no infra; respx for LLM HTTP; includestest_bot_boundary.py(AST check of the hexagonal bot rule). - Integration (
tests/integration/, 21 files / 134 tests,-m integration): full docker-compose test stack, alembic upgrade, seeded service token, ASGI httpx + LifespanManager; covers upload pipeline, each worker, credits/quota DB, auth flows, passkeys/magic links, B2B, admin panel, billing checkout/webhook/pay page, subscriptions, refunds. tests/conftest.pydeduplicatessrc.contract_check.*vscontract_check.*module loading (avoids duplicate Prometheus timeseries).
12. Cross-cutting conventions
- Money is integer kopecks everywhere; decimals exist only at the ЮKassa HTTP boundary.
- Enum-ish columns are
TEXT + CHECK; migrations additive and hand-written. - Repositories never commit; callers own transactions.
- Every queue handler is idempotent via
documents.statusre-read; refunds/slot releases are exactly-once via guards. - Publisher confirms on every publish that follows a paid reservation.
- External providers (LLM, payments, storage) sit behind ports in
core/; adapters are swappable and respx-testable. - Feature degrade switches:
PLANS_ENABLED,YOOKASSA_ENABLED,PRESCREEN_ENABLED,PRESCREEN_LLM_FALLBACK_ENABLED,SMTP_HOST(empty = dev logger),REFUND_POLICY. - Domain language (Plan/Subscription/Quota/Credits/Topup/Invoice/Renewal/
Saved Payment Method/Dunning/Signup Bonus/ЕРИП/Billing Hold/Clawback/
Document Slot) is normative — see root
CONTEXT.md; ADRs live indocs/adr/.