Skip to content

About

Async Telegram bot built with aiogram for LinkTracker: manages GitHub and Stack Overflow subscriptions via the LinkTracker Scrapper API and receives update notifications over Kafka or HTTP.

Topics

Resources

Stars

0 stars

Watchers

0 watching

Forks

Repository files navigation

LinkTracker Bot Service

CI Python License: MIT Ruff Code style: black Checked with mypy

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.


Requirements

  • Python 3.12+
  • Poetry
  • GNU Make (optional, see available targets below)

Setup

poetry install

Create a .env file in the project root:

cp .env.example .env

Fill 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

Run

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 -d

That 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-updates topic is created with the expected settings.

Then start the bot from this repo:

poetry run python src/main.py

To 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.


Tests and linting

make help    # Show available targets
make test    # Run pytest
make lint    # Run black, ruff, and dead fixture checks

You can also run tools directly via Poetry:

poetry run pytest ./tests
poetry run black ./src ./tests
poetry run ruff check ./src ./tests

CI 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.


Supported commands

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

Incoming updates from Scrapper

The bot receives link-update notifications from Scrapper via one of two transports, selected by the UPDATES_TRANSPORT environment variable:

HTTP

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.

Kafka (default)

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.

Switching to the update pipeline processed-updates chain

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.

Message format

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 via SCHEMA_REGISTRY_URL (default http://localhost:8081).
  • json — UTF-8 JSON matching the same LinkUpdate schema.

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.

Retry and Dead Letter Queue

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).

DLQ topic settings

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).

Avro contract

When KAFKA_VALUE_FORMAT=avro, both services agree on:

  • schema fullname linktracker.LinkUpdate (record LinkUpdate, namespace linktracker);
  • fields id: long, url: string, description: string, tgChatIds: array<long>;
  • subject link-updates-value (Confluent TopicNameStrategy);
  • 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.

Verifying the Scrapper → Kafka → Bot integration locally

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_path registers the Avro schema in Schema Registry, encodes a message in the Confluent wire format with a vanilla aiokafka producer (emulating the Scrapper publisher), and consumes it through the real LinkUpdatesConsumer + AvroPayloadDecoder, asserting the MessageSender receives the expected calls:

    poetry run pytest tests/test_avro_kafka_integration.py::test_avro_full_path -v
  • Scrapper side — tests/test_avro_integration.py exercises 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.


Logging

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

Reliability

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.

Bot → Scrapper

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, with constant (default) or exponential backoff scaled by HTTP_RETRY_BACKOFF_SECONDS. Only statuses listed in HTTP_RETRY_STATUSES (default 500,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.CircuitBreaker with error_ratio = CB_FAILURE_RATE_THRESHOLD, sliding window CB_SLIDING_WINDOW_SECONDS, and broken/recovery durations CB_WAIT_DURATION_IN_OPEN_SECONDS / CB_RECOVERY_TIME_SECONDS. Disable with CB_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 → Telegram

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.

Rate limit on POST /updates

The public ingress endpoint is rate-limited per client IP via aiohttp-ratelimiter:

  • RATE_LIMIT_DEFAULT is parsed by the limits library; the format is <count>/<unit> with units second, minute, hour, day (default 60/minute).
  • Disable with RATE_LIMIT_ENABLED=false (no decorator is attached).
  • On exceedance the server returns 429 with an ApiErrorResponse body (code="429", exceptionName="RateLimitExceeded").

Observability

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 bot

See OBSERVABILITY.md for the metric contract, label conventions, PromQL examples, and the importable Grafana dashboard at observability/grafana/dashboards/bot.json.


Project structure

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

License

Licensed under the MIT License. See LICENSE.

About

Async Telegram bot built with aiogram for LinkTracker: manages GitHub and Stack Overflow subscriptions via the LinkTracker Scrapper API and receives update notifications over Kafka or HTTP.

Topics

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages