Telegram bot for the LinkTracker project.
Handles user commands (/start, /help, /track, /untrack, /list) and receives
link-update notifications from the Scrapper service via HTTP (POST /updates) or Kafka consumer.
Supports both HTTP and gRPC transports for outgoing communication with Scrapper.
The backend Scrapper service is maintained in a separate repository: linktracker-scrapper.
- Python 3.12+
- Poetry
- GNU Make (optional, see available targets below)
poetry installCreate a .env file in the project root:
cp .env.example .envFill in the required values (get a bot token from BotFather):
BOT_TOKEN=telegram-bot-token-here
SCRAPPER_URL=http://localhost:8080
# Optional (defaults shown):
# SCRAPPER_TRANSPORT=http # http or grpc
# SCRAPPER_GRPC_TARGET=localhost:50051
# SERVER_HOST=0.0.0.0
# SERVER_PORT=8000
# UPDATES_TRANSPORT=kafka # kafka (default) or http
# KAFKA_BOOTSTRAP_SERVERS=localhost:29092,localhost:29093,localhost:29094
# KAFKA_TOPIC=link-updates
# KAFKA_CONSUMER_GROUP=bot-service
# KAFKA_DLQ_TOPIC=link-updates-dlq
# KAFKA_RETRY_ATTEMPTS=3 # processing-error retries before DLQ
# KAFKA_RETRY_DELAY_MS=0 # fixed delay between retry attempts
# KAFKA_VALUE_FORMAT=avro # avro (default) or json
# SCHEMA_REGISTRY_URL=http://localhost:8081
# Outgoing HTTP resilience (Bot -> Scrapper, also session timeout for Bot -> Telegram):
# HTTP_TIMEOUT_SECONDS=10.0
# HTTP_RETRY_MAX_ATTEMPTS=3
# HTTP_RETRY_BACKOFF_SECONDS=0.5
# HTTP_RETRY_STRATEGY=constant # constant or exponential
# HTTP_RETRY_STATUSES=500,502,503,504
# Circuit breaker for Bot -> Scrapper:
# CB_ENABLED=true
# CB_FAILURE_RATE_THRESHOLD=0.5
# CB_SLIDING_WINDOW_SECONDS=10.0
# CB_WAIT_DURATION_IN_OPEN_SECONDS=30.0
# CB_RECOVERY_TIME_SECONDS=5.0
# Per-IP rate limit on POST /updates:
# RATE_LIMIT_ENABLED=true
# RATE_LIMIT_DEFAULT=60/minute
The bot expects an external Kafka cluster and Schema Registry (Avro is the default wire format). The canonical local stack lives in the Scrapper repo — clone it next to this project and start the infrastructure from there:
# In the Scrapper repo:
docker compose up -dThat brings up:
- a 3-broker Kafka cluster in KRaft mode, exposed on the host as
localhost:29092,localhost:29093,localhost:29094; - Schema Registry on
http://localhost:8081.
These addresses match this project's defaults (KAFKA_BOOTSTRAP_SERVERS,
SCHEMA_REGISTRY_URL), so the bot needs no extra wiring once the stack is
up. The Scrapper service registers the Avro schema and creates the
link-updates topic on its own startup; the bot creates link-updates-dlq
on its startup.
Note: before the first full bot run, bring up the infrastructure stack and start Scrapper at least once so that the
link-updatestopic is created with the expected settings.
Then start the bot from this repo:
poetry run python src/main.pyTo switch to the JSON wire format (e.g. when running without Schema
Registry) set KAFKA_VALUE_FORMAT=json in .env. To use the legacy HTTP
ingress instead of Kafka, set UPDATES_TRANSPORT=http.
make help # Show available targets
make test # Run pytest
make lint # Run black, ruff, and dead fixture checksYou can also run tools directly via Poetry:
poetry run pytest ./tests
poetry run black ./src ./tests
poetry run ruff check ./src ./testsCI runs lint, unit tests, and Kafka integration tests. The Telegram e2e tests
(tests/test_e2e.py) require a real Telegram bot token and outbound network
access, so they are intended for manual runs and are excluded from the default
CI.
| Command | Description |
|---|---|
/start |
Greets the user by name; registers the chat in Scrapper |
/help |
Shows a list of available commands |
/track |
Start tracking a link (dialog) |
/untrack |
Stop tracking a link (dialog) |
/list [tag] |
Show tracked links, optionally by tag |
/cancel |
Cancel current dialog |
| Unknown | Replies with an error and suggests /help |
The bot receives link-update notifications from Scrapper via one of two transports,
selected by the UPDATES_TRANSPORT environment variable:
The bot exposes an HTTP server (default 0.0.0.0:8000):
| Method | Endpoint | Description |
|---|---|---|
| POST | /updates |
Receive link update from Scrapper |
Request body must match the LinkUpdate schema. Returns 200 OK on success,
400 Bad Request with ApiErrorResponse on invalid input.
See the OpenAPI contract for full schema details.
The bot consumes messages from a Kafka topic (link-updates by default).
Each message represents a LinkUpdate record (id, url, description,
tgChatIds); the description field is sent to Telegram as-is. Only id
and tgChatIds are mandatory; url and description are optional and any
extra fields are ignored.
The full pipeline is Scrapper -> link.raw-updates -> update pipeline -> link.processed-updates -> Bot. To consume the update pipeline output instead of
the raw Scrapper topic, set in .env:
KAFKA_TOPIC=link.processed-updates
KAFKA_VALUE_FORMAT=json
The processed payload is {id, description, tgChatIds, priority}; url
is omitted and priority is currently informational (no Bot-side logic).
Offsets are committed per message (enable_auto_commit=False) only after the
message is either successfully handled or reliably forwarded to the DLQ.
The wire format is selected by KAFKA_VALUE_FORMAT:
avro(default) — Confluent wire format:[0x00][schema_id: 4 bytes BE][avro body]. The schema is fetched from Schema Registry by id (subject<topic>-value) and cached in memory. Configure the registry endpoint viaSCHEMA_REGISTRY_URL(defaulthttp://localhost:8081).json— UTF-8 JSON matching the sameLinkUpdateschema.
A malformed payload (bad magic byte, broken Avro body, invalid JSON) is
routed to the DLQ as decode_error without retry.
Schema Registry-related failures are never sent to the DLQ:
- transport errors, timeouts, and 5xx responses are retried internally by the registry client up to N times;
- 4xx responses, unparseable schema documents, and exhausted retries bubble up to the consumer.
In all of these cases the consumer stops without committing the offset, so the message is not acknowledged and can be reprocessed after the underlying problem is fixed.
Consumer errors are split into three categories:
| Error category | Retry? | DLQ | Header reason |
|---|---|---|---|
| Decode error | no | immediate | decode_error |
| Validation error | no | immediate | validation_error |
| Processing error | up to N attempts | on exhaust | processing_error |
The number of attempts is controlled by KAFKA_RETRY_ATTEMPTS (default 3);
KAFKA_RETRY_DELAY_MS adds an optional fixed delay between attempts.
DLQ messages carry the unmodified original payload as the Kafka record value.
Metadata is attached as Kafka headers: reason, error (short message),
attempts (processing-error attempts, 0 for decode/validation). The target
topic is configured via KAFKA_DLQ_TOPIC (default link-updates-dlq).
The bot creates the DLQ topic on startup (idempotently, via
AIOKafkaAdminClient.create_topics) with the following settings:
| Setting | Value | Reason |
|---|---|---|
num_partitions |
3 |
one partition per broker so manual replay can be parallelised |
replication_factor |
3 |
matches the durability profile of link-updates; survives one broker |
min.insync.replicas |
2 |
with acks=all no DLQ writes are lost when one broker is down |
retention.ms |
30 days |
long enough for an operator to triage and replay failed messages |
The main link-updates topic is created and configured by the Scrapper
service (see its README for the rationale).
When KAFKA_VALUE_FORMAT=avro, both services agree on:
- schema fullname
linktracker.LinkUpdate(recordLinkUpdate, namespacelinktracker); - fields
id: long,url: string,description: string,tgChatIds: array<long>; - subject
link-updates-value(ConfluentTopicNameStrategy); - compatibility level
BACKWARD(Schema Registry default).
Scrapper registers the schema at startup and embeds the resulting
schema_id in every message; the bot fetches the schema by id on first
use and caches it in memory.
The end-to-end path is split across two repos. Both tests use Testcontainers to boot real Kafka and Schema Registry instances, so Docker must be running.
-
Bot side —
tests/test_avro_kafka_integration.py::test_avro_full_pathregisters the Avro schema in Schema Registry, encodes a message in the Confluent wire format with a vanillaaiokafkaproducer (emulating the Scrapper publisher), and consumes it through the realLinkUpdatesConsumer+AvroPayloadDecoder, asserting theMessageSenderreceives the expected calls:poetry run pytest tests/test_avro_kafka_integration.py::test_avro_full_path -v
-
Scrapper side —
tests/test_avro_integration.pyexercises the same contract from the producer end:UpdateChecker→AvroKafkaNotificationSender→ topic, with the message read back raw and validated against the contract. Run it from the Scrapper repo.
Together the two tests prove the wire-format contract end-to-end.
The bot uses structured logging in logfmt format. Extra fields are appended as key=value pairs after a | separator.
Values containing spaces are automatically quoted.
Example output:
2026-02-21 12:00:00,000 INFO __main__ bot_started | polling=long_polling
2026-02-21 12:00:01,000 INFO bot.handlers.common command_received | command=/start user_id=123456789
2026-02-21 12:00:02,000 INFO bot.handlers.common command_received | command=/help user_id=123456789
2026-02-21 12:00:03,000 INFO bot.handlers.common unknown_command | command=/unknown user_id=123456789
Outbound HTTP and the public POST /updates endpoint are wrapped with
configurable resilience primitives. All thresholds are exposed via
pydantic-settings, no values are hard-coded.
HttpScrapperClient routes every request through ResilientHttpCaller,
which composes CircuitBreaker(Retry(aiohttp request with timeout)):
- Timeout —
HTTP_TIMEOUT_SECONDS(aiohttp.ClientTimeout(total=...)). - Retry — up to
HTTP_RETRY_MAX_ATTEMPTS, withconstant(default) orexponentialbackoff scaled byHTTP_RETRY_BACKOFF_SECONDS. Only statuses listed inHTTP_RETRY_STATUSES(default500,502,503,504) and transient network errors (aiohttp.ClientError,asyncio.TimeoutError) are retried. 4xx responses are returned to the caller without retry. - Circuit breaker —
aiomisc.CircuitBreakerwitherror_ratio = CB_FAILURE_RATE_THRESHOLD, sliding windowCB_SLIDING_WINDOW_SECONDS, and broken/recovery durationsCB_WAIT_DURATION_IN_OPEN_SECONDS/CB_RECOVERY_TIME_SECONDS. Disable withCB_ENABLED=false.
The breaker observes the outcome of fully exhausted retry sequences, so
a single transient failure that recovers on retry does not count toward
the breaker threshold. After retry exhaustion every transient failure
surfaces as TransientHttpError; an open breaker surfaces as
CircuitOpenError (subclass of TransientHttpError).
Bot is constructed with AiohttpSession(timeout=HTTP_TIMEOUT_SECONDS),
which bounds Telegram API call duration. Retries and a circuit breaker
are deliberately not added here: Telegram sendMessage is not
idempotent, so naive retries on transient failures would risk delivering
duplicate notifications to users. If retries are ever needed, they
should be designed alongside an idempotency strategy, for example
deduplication by a stable notification id per chat.
The public ingress endpoint is rate-limited per client IP via
aiohttp-ratelimiter:
RATE_LIMIT_DEFAULTis parsed by thelimitslibrary; the format is<count>/<unit>with unitssecond,minute,hour,day(default60/minute).- Disable with
RATE_LIMIT_ENABLED=false(no decorator is attached). - On exceedance the server returns
429with anApiErrorResponsebody (code="429",exceptionName="RateLimitExceeded").
The bot exposes Prometheus metrics on /metrics at port 8011 (configurable
via METRICS_HOST / METRICS_PORT). The endpoint is bound inside the
bot-service container only and is not published to the host.
The shared Prometheus/Grafana stack is provisioned by the Scrapper repository.
The bot joins its external Docker network so Prometheus scrapes
bot-service:8011:
# In the Scrapper repo:
docker compose up -d prometheus grafana
# In this repo:
docker compose up -d botSee OBSERVABILITY.md for the metric contract, label
conventions, PromQL examples, and the importable Grafana dashboard at
observability/grafana/dashboards/bot.json.
src/
├── main.py # Entrypoint, composition root
├── config.py # Configuration (pydantic-settings)
├── logging_config.py # Structured logging setup
├── application/
│ ├── commands.py # Command registry
│ ├── models.py # Pydantic models (OpenAPI contract)
│ ├── notification_handler.py # Link-update processing logic
│ ├── clients/
│ │ ├── scrapper.py # ScrapperClient protocol
│ │ ├── message_sender.py # MessageSender protocol
│ │ └── errors.py # Client error types
│ └── use_cases/ # Business logic (pure Python)
│ ├── start.py
│ ├── help.py
│ ├── unknown.py
│ ├── track.py
│ ├── untrack.py
│ └── list_links.py
├── bot/
│ ├── commands.py # Telegram menu commands
│ ├── states.py # FSM state groups
│ ├── server.py # HTTP server (POST /updates)
│ ├── message_sender.py # MessageSender implementation (aiogram)
│ └── handlers/
│ ├── common.py # /start, /help, unknown
│ ├── track.py # /track dialog (FSM)
│ ├── untrack.py # /untrack dialog (FSM)
│ └── list_links.py # /list command
├── infrastructure/
│ ├── scrapper_client.py # HTTP ScrapperClient via ResilientHttpCaller
│ ├── resilient_http_caller.py # Timeout/retry/circuit breaker wrapper for outgoing HTTP
│ ├── grpc_scrapper_client.py # gRPC implementation of ScrapperClient
│ ├── kafka_consumer.py # Kafka consumer with retry and DLQ dispatch
│ ├── kafka_topic_creator.py # Idempotent DLQ topic creation on startup
│ ├── dlq_producer.py # Kafka Dead Letter Queue producer
│ ├── payload_decoder.py # PayloadDecoder protocol + JSON decoder
│ ├── avro_payload_decoder.py # Confluent-wire-format Avro decoder
│ └── schema_registry_client.py # Async Schema Registry REST client
├── scrapper_pb2.py # Generated protobuf code
└── scrapper_pb2_grpc.py # Generated gRPC stubs
proto/
└── scrapper.proto # Protobuf service definition
Licensed under the MIT License. See LICENSE.