diff --git a/homework/05-message-queue/.dockerignore b/homework/05-message-queue/.dockerignore new file mode 100644 index 0000000..d291891 --- /dev/null +++ b/homework/05-message-queue/.dockerignore @@ -0,0 +1,5 @@ +target +**/target +**/__pycache__ +**/.pytest_cache +*.pyc diff --git a/homework/05-message-queue/docker-compose.yml b/homework/05-message-queue/docker-compose.yml new file mode 100644 index 0000000..df61c30 --- /dev/null +++ b/homework/05-message-queue/docker-compose.yml @@ -0,0 +1,48 @@ +name: distsys-mq +services: + rabbitmq: + image: rabbitmq:4.3.6-management-alpine + volumes: + - ./tests/rabbitmq.conf:/etc/rabbitmq/rabbitmq.conf + ports: + - 15672:15672 + - 5672:5672 + + server: + image: distsys-mq-server + build: + context: ./solution/server + volumes: + - data-volume:/data + ports: + - 5000:5000 + + server-fdv: + image: distsys-mq-server + build: + context: ./solution/server + volumes: + - fake-data-volume:/data + ports: + - 5000:5000 + + worker: + image: distsys-mq-worker + build: + context: ./solution/worker + volumes: + - data-volume:/data + deploy: + replicas: 2 + + network-fault: + image: distsys-mq-network-fault + build: + context: ./tests/network-fault + network_mode: service:server + cap_add: + - NET_ADMIN + +volumes: + data-volume: + fake-data-volume: diff --git a/homework/05-message-queue/media/architecture.png b/homework/05-message-queue/media/architecture.png new file mode 100644 index 0000000..7afa5be Binary files /dev/null and b/homework/05-message-queue/media/architecture.png differ diff --git a/homework/05-message-queue/media/architecture.svg b/homework/05-message-queue/media/architecture.svg new file mode 100644 index 0000000..607a867 --- /dev/null +++ b/homework/05-message-queue/media/architecture.svg @@ -0,0 +1,74 @@ + + Архитектура сервиса генерации описаний + Пользователь обращается к REST API. Сервер передаёт задачи через RabbitMQ воркерам. Воркеры сохраняют описания в общем Docker volume, откуда сервер читает их. Для компонента CALLBACKS воркеры также отправляют уведомления серверу через RabbitMQ. + + + + + + + + + + + + + + + + + + Пользователь + HTTP-клиент + + + Сервер + REST API + + + + HTTP + + + + RabbitMQ + + Задачи + + + Уведомления + + + + Воркеры + + Воркер 1 + + Воркер 2 + + + + + + + + + Общий Docker volume + Файлы с описаниями + + + Чтение описаний + + Запись описаний + + + Уведомления для компонента CALLBACKS + diff --git a/homework/05-message-queue/readme.md b/homework/05-message-queue/readme.md new file mode 100644 index 0000000..c088e83 --- /dev/null +++ b/homework/05-message-queue/readme.md @@ -0,0 +1,152 @@ +# Практика с RabbitMQ + +В этом задании вы научитесь работать с брокером сообщений [RabbitMQ](https://rabbitmq.com/getstarted.html). Вам предстоит реализовать сервис, который принимает URL изображения, асинхронно генерирует описание, сохраняет его в файл и возвращает по запросу. Генерация описания имитируется готовой функцией из заготовки. + +## Архитектура и интерфейс сервиса + +Очередь сообщений позволяет передать длительную задачу на асинхронную обработку. Она помогает распределять работу между несколькими воркерами и сглаживать всплески нагрузки. В этом задании вам также предстоит обеспечить обработку задач при временных отказах. + +Архитектура сервиса + +1. Сервер с REST API принимает запросы пользователей и передаёт задачи на генерацию описаний через RabbitMQ. +2. Воркеры получают задачи, генерируют описания и записывают их в общее хранилище, откуда сервер читает результаты. В задании хранилище — это [Docker volume](https://docs.docker.com/engine/storage/volumes/), см. [docker-compose.yml](docker-compose.yml). + +Версия RabbitMQ для проверки задана в [docker-compose.yml](docker-compose.yml). Используйте это окружение для локального тестирования. + +### REST API + +Тела запросов и успешных ответов передаются в JSON. Успешный запрос возвращает код `200`; формат тела ответа при ошибке не регламентируется. Идентификатор изображения может быть строкой или целым числом; в ответах POST и GET он должен иметь одинаковый тип. + +``` +POST /api/v1.0/images +Принимает URL изображения для обработки и возвращает id изображения. +Каждый запрос получает новый уникальный id, даже при повторной отправке того же URL. +Если JSON некорректен или image_url отсутствует, пуст или не является строкой, +верните код 400. + +Тело запроса: +{ + "image_url": str +} + +Тело ответа: +{ + "image_id": str | int +} +``` + +``` +GET /api/v1.0/images +Возвращает id всех обработанных изображений без повторов. Порядок не важен. + +Тело ответа: +{ + "image_ids": List[str | int] +} +``` + +``` +GET /api/v1.0/images/ +Возвращает описание для данного изображения, если оно уже было обработано, +и код 404, если id неизвестен или обработка ещё не завершена. + +Тело ответа: +{ + "caption": str +} +``` + +## Компоненты задания и оценивание + +За выполнение задания можно получить 10 баллов: + +| Группа тестов | Баллы | Что проверяется | +| --- | ---: | --- | +| `BASIC_API` | 1 | Пустой список результатов, отклонение некорректных запросов и ответ `404` для неизвестного ID. | +| `BASIC_PROCESSING` | 3 | Передача задач через RabbitMQ, обработка воркерами и получение результатов через REST API без отказов. | +| `CALLBACKS` | 2 | Уведомления о завершении обработки через RabbitMQ: сервер отвечает на `GET .../images` без чтения списка файлов в общей директории. | +| `FAULT_TOLERANCE_1` | 1 | Работа сервиса после простоя соединений с брокером. | +| `FAULT_TOLERANCE_2` | 1 | Сохранность запросов при недоступности и перезапуске брокера во время отправки задач. | +| `FAULT_TOLERANCE_3` | 1 | Завершение задач после отказа одного или обоих воркеров. | +| `FAULT_TOLERANCE_4` | 1 | Восстановление после отказов брокера и воркеров, сохранность уведомлений и обработка повторных доставок. | + +Баллы за каждую группу начисляются при прохождении всех её тестов. Провал одной группы не обнуляет баллы за остальные. Группа `FAULT_TOLERANCE_4` предполагает реализацию уведомлений через RabbitMQ; передавать в них сами описания изображений не требуется. + +При проверке отказоустойчивости уже запущенный REST API должен принимать корректные запросы с кодом `200` и возвращать готовые результаты, даже когда брокер или воркеры недоступны. Запрос считается принятым, если POST вернул `200`. После восстановления брокера, связи и хотя бы одного воркера каждый принятый запрос должен быть обработан, а его описание — доступно через API. Таймаут каждого HTTP-запроса в тестах — 3 секунды. + +Сервер и общее хранилище в этих сценариях не отказывают: восстановление после их перезапуска или потери данных не требуется. Повторная доставка задач и уведомлений возможна. Повторная обработка задачи допустима, но её ID должен появляться в списке результатов только один раз. + +Приложите отчёт `solution/readme.md` с описанием устройства решения и выполненных компонентов. Для отказоустойчивости обоснуйте, почему принятые запросы не теряются. На защите нужно объяснить работу реализованных компонентов и подтвердить это обоснование. Общие требования к отчёту приведены в [правилах курса](../readme.md#отчёт). + +## Заготовки для решения + +В папке `solution/server` находится заготовка для сервера с реализацией REST API на Flask. + +В папке `solution/worker` содержится заготовка для воркера. Для генерации описания используйте функцию `produce_image_caption`, передавая ей строку `image_url`. Скачивать и сохранять само изображение не нужно: функция имитирует его обработку. В файл сохраняется только полученное описание в кодировке UTF-8. + +Очередь задач должна называться `task_queue`, а описания должны сохраняться в файлы `/data/{image_id}.txt`. Эти соглашения используются в тестах. + +Для взаимодействия с RabbitMQ на Python предлагается использовать библиотеку [Pika](https://github.com/pika/pika) (см. семинар 5). + +Весь код решения должен размещаться в папке `solution`. При сдаче решения в тестирующую систему отправляется только эта папка, изменения вне неё учитываться не будут. + +## Порядок выполнения задания + +Начните с [материалов семинара 5](../../materials/05-indirect-comm/seminar/) и примера [work_queues](../../materials/05-indirect-comm/seminar/work_queues/). Дополнительно можно обращаться к официальным [туториалам RabbitMQ](https://www.rabbitmq.com/tutorials). + +1. Реализуйте передачу задач от REST API через RabbitMQ, обработку воркерами и получение результатов. Проверьте группы `BASIC_API` и `BASIC_PROCESSING`. При необходимости измените Dockerfile сервера и воркера. +2. Добавьте уведомления сервера о завершении обработки и проверьте группу `CALLBACKS`. Полезный пример обмена запросами и ответами есть в [туториале по RPC](https://www.rabbitmq.com/tutorials/tutorial-six-python). При использовании Pika учитывайте [ограничения потокобезопасности](https://pika.github.io/pika/latest/faq/). +3. Работу над отказоустойчивостью начните с `test_heartbeats_timeout` (`FAULT_TOLERANCE_1`). Изучите [документацию по heartbeat](https://www.rabbitmq.com/docs/heartbeats) и настройки в [rabbitmq.conf](tests/rabbitmq.conf). Для `BlockingConnection` полезно [обсуждение работы соединения при редкой отправке сообщений](https://github.com/pika/pika/discussions/1382). +4. Подумайте, в каких случаях ваше решение может потерять принятый запрос. Изучите [подтверждения доставки в RabbitMQ](https://www.rabbitmq.com/docs/confirms); для Pika также полезен [пример асинхронного отправителя](https://github.com/pika/pika/blob/main/examples/asynchronous_publisher_example.py). Проверьте отдельно отказы брокера при отправке задач (`FAULT_TOLERANCE_2`) и отказы воркеров (`FAULT_TOLERANCE_3`). +5. Перейдите к комплексным сценариям `FAULT_TOLERANCE_4`: перезапуску брокера, сохранности уведомлений и повторной доставке. Сопоставьте поведение решения с [руководством RabbitMQ по надёжности](https://www.rabbitmq.com/docs/reliability). Подумайте, какие моменты отказа ещё стоит проверить самостоятельно. +6. Подготовьте отчёт с обоснованием отказоустойчивости реализованных компонентов, запустите все тесты и сдайте решение с отчётом в тестирующую систему. + +При отладке сопоставляйте логи компонентов с состоянием очередей в веб-интерфейсе RabbitMQ: `http://localhost:15672`, логин и пароль — `guest`. Как на семинаре, обращайте внимание на сообщения Ready и Unacked. + +## Тестирование решения + +Тесты, проверяющие решение, находятся в папке `tests`. + +Бонусы за пробелы в тестах начисляются по [общим правилам](../readme.md#бонусы-за-пробелы-в-тестах). + +### Локальное тестирование + +Команды ниже выполняйте из папки задания. Для локального запуска тестов установите зависимости из `tests/requirements.txt`. + +Перед запуском тестов соберите образы: + +``` +docker compose build +``` + +Запуск всех тестов выполняется с помощью команды: + +``` +python3 tests/main.py +``` + +Отдельный тест можно запустить так: + +``` +pytest -vs --tb=short tests/test_server.py::test_single_image +``` + +Для тестирования в окружении, аналогичном тестирующей системе, выполните из папки задания: + +```bash +docker run --privileged --pull always --rm -v ./solution:/hw/solution distsys.ru/course/message-queue:latest +``` + +### Проверка в тестирующей системе + +Отправьте ваше решение в тестирующую систему следуя [инструкции](../readme.md) и дождитесь результатов. + +## ЧаВо + +**Можно ли реализовать решение не на Python?** + +Да. Замените заготовки и Dockerfile сервера и воркера. Реализуйте аналог функции-заглушки, возвращающий строку по `image_url`; точное совпадение результата с Python-версией не требуется. + +**Можно ли обрабатывать изображения локально на сервере?** + +Нет. Описания должны генерироваться только воркерами. Решение, которое генерирует их на сервере, не засчитывается (0 баллов). diff --git a/homework/05-message-queue/solution/server/Dockerfile b/homework/05-message-queue/solution/server/Dockerfile new file mode 100644 index 0000000..27a1a1b --- /dev/null +++ b/homework/05-message-queue/solution/server/Dockerfile @@ -0,0 +1,13 @@ +# syntax=docker/dockerfile:1 +FROM python:3.12-slim + +RUN apt-get update && apt-get install -y curl \ + && rm -rf /var/lib/apt/lists/* + +WORKDIR /server +COPY requirements.txt . +RUN --mount=type=cache,id=distsys-course-pip,target=/root/.cache/pip,sharing=locked \ + pip install -r requirements.txt +COPY . . + +CMD ["python3", "-u", "server.py"] diff --git a/homework/05-message-queue/solution/server/requirements.txt b/homework/05-message-queue/solution/server/requirements.txt new file mode 100644 index 0000000..c772219 --- /dev/null +++ b/homework/05-message-queue/solution/server/requirements.txt @@ -0,0 +1,2 @@ +flask==3.1.2 +pika==1.3.2 \ No newline at end of file diff --git a/homework/05-message-queue/solution/server/server.py b/homework/05-message-queue/solution/server/server.py new file mode 100644 index 0000000..95b9886 --- /dev/null +++ b/homework/05-message-queue/solution/server/server.py @@ -0,0 +1,51 @@ +from flask import Flask, request +from typing import List, Optional + + +class Server: + # TODO: Implement API server + def __init__(self, mq_host, mq_port, data_dir): + pass + + def add_image(self, image_url: str) -> str: + raise NotImplementedError + + def get_processed_images(self) -> List[str]: + raise NotImplementedError + + def get_image_caption(self, image_id: str) -> Optional[str]: + raise NotImplementedError + + +def create_app() -> Flask: + app = Flask(__name__) + + server = Server('rabbitmq', 5672, '/data') + + @app.route('/api/v1.0/images', methods=['POST']) + def add_image(): + body = request.get_json(force=True) + if not isinstance(body, dict) or not isinstance(body.get('image_url'), str) or not body['image_url']: + return 'image_url must be a nonempty string', 400 + image_id = server.add_image(body['image_url']) + return {"image_id": image_id} + + @app.route('/api/v1.0/images', methods=['GET']) + def get_processed_images(): + image_ids = server.get_processed_images() + return {"image_ids": image_ids} + + @app.route('/api/v1.0/images/', methods=['GET']) + def get_image_caption(image_id): + result = server.get_image_caption(image_id) + if result is None: + return "Image not found.", 404 + else: + return {'caption': result} + + return app + + +if __name__ == '__main__': + app = create_app() + app.run(host='0.0.0.0', port=5000) diff --git a/homework/05-message-queue/solution/worker/Dockerfile b/homework/05-message-queue/solution/worker/Dockerfile new file mode 100644 index 0000000..93ed6fd --- /dev/null +++ b/homework/05-message-queue/solution/worker/Dockerfile @@ -0,0 +1,13 @@ +# syntax=docker/dockerfile:1 +FROM python:3.12-slim + +RUN apt-get update && apt-get install -y curl \ + && rm -rf /var/lib/apt/lists/* + +WORKDIR /worker +COPY requirements.txt . +RUN --mount=type=cache,id=distsys-course-pip,target=/root/.cache/pip,sharing=locked \ + pip install -r requirements.txt +COPY . . + +CMD ["python3", "-u", "worker.py"] diff --git a/homework/05-message-queue/solution/worker/requirements.txt b/homework/05-message-queue/solution/worker/requirements.txt new file mode 100644 index 0000000..cde8834 --- /dev/null +++ b/homework/05-message-queue/solution/worker/requirements.txt @@ -0,0 +1 @@ +pika==1.3.2 \ No newline at end of file diff --git a/homework/05-message-queue/solution/worker/worker.py b/homework/05-message-queue/solution/worker/worker.py new file mode 100644 index 0000000..7ad9eca --- /dev/null +++ b/homework/05-message-queue/solution/worker/worker.py @@ -0,0 +1,15 @@ +class Worker: + # TODO: Implement server + def __init__(self, mq_host, mq_port, data_dir): + pass + + def produce_image_caption(self, image_url): + return str(abs(hash(image_url)) % (10 ** 8)) + + def run(self): + raise NotImplementedError + + +if __name__ == '__main__': + worker = Worker('rabbitmq', 5672, '/data') + worker.run() diff --git a/homework/05-message-queue/tests/main.py b/homework/05-message-queue/tests/main.py new file mode 100644 index 0000000..b41112a --- /dev/null +++ b/homework/05-message-queue/tests/main.py @@ -0,0 +1,107 @@ +import argparse +import pathlib +import pytest + +from collections import defaultdict + +SCRIPT_DIR = pathlib.Path(__file__).parent.resolve() + +TEST_GROUPS = { + 'BASIC_API': { + 'tests': [ + 'test_empty_data_dir', + 'test_bad_request', + 'test_nonexistent_image', + ], + 'points': 1 + }, + 'BASIC_PROCESSING': { + 'tests': [ + 'test_task_queue', + 'test_single_image', + 'test_multiple_images', + 'test_captions_generated_on_workers' + ], + 'points': 3 + }, + 'CALLBACKS': { + 'tests': [ + 'test_multiple_images_no_listdir' + ], + 'points': 2 + }, + 'FAULT_TOLERANCE_1': { + 'tests': [ + 'test_heartbeats_timeout' + ], + 'points': 1 + }, + 'FAULT_TOLERANCE_2': { + 'tests': [ + 'test_publisher_confirms' + ], + 'points': 1 + }, + 'FAULT_TOLERANCE_3': { + 'tests': [ + 'test_faulty_worker', + 'test_two_faulty_workers', + ], + 'points': 1 + }, + 'FAULT_TOLERANCE_4': { + 'tests': [ + 'test_faulty_worker_and_rabbit_restart', + 'test_total_eclipse_of_the_heart', + 'test_notifications_survive_rabbit_restart', + 'test_duplicate_deliveries' + ], + 'points': 1 + } +} + +class PassedCounter: + def __init__(self): + self.test_to_group = {} + for group_name, group in TEST_GROUPS.items(): + for test in group['tests']: + self.test_to_group[test] = group_name + self.passed_by_group = defaultdict(int) + + def pytest_report_teststatus(self, report, config): + if report.when == 'call' and report.passed: + test = report.nodeid.split('::')[1].split('[')[0] + group_name = self.test_to_group[test] + self.passed_by_group[group_name] += 1 + + +def main(argv=None): + parser = argparse.ArgumentParser() + parser.add_argument('--ci', action='store_true', + help='Fail unless all reference solution tests pass') + args = parser.parse_args(argv) + counter = PassedCounter() + test_status = pytest.main( + ['-vs', '--tb=short', str(SCRIPT_DIR / 'test_server.py')], plugins=[counter]) + + score = 0 + print() + for group_name, group in TEST_GROUPS.items(): + total = len(group['tests']) + passed = counter.passed_by_group[group_name] + print(f'Test group {group_name}: passed {passed} of {total} tests') + if passed == total: + score += group['points'] + + print(f"\nSCORE: {score}") + if args.ci: + max_score = sum(group['points'] for group in TEST_GROUPS.values()) + return int(test_status) or int(score != max_score) + # A failed student test still produces a valid partial score. + if test_status == pytest.ExitCode.TESTS_FAILED: + return 0 + return int(test_status) + + +if __name__ == '__main__': + raise SystemExit(main()) diff --git a/homework/05-message-queue/tests/network-fault/Dockerfile b/homework/05-message-queue/tests/network-fault/Dockerfile new file mode 100644 index 0000000..6f18b90 --- /dev/null +++ b/homework/05-message-queue/tests/network-fault/Dockerfile @@ -0,0 +1,4 @@ +FROM docker:dind + +ENTRYPOINT ["sleep"] +CMD ["infinity"] diff --git a/homework/05-message-queue/tests/rabbitmq.conf b/homework/05-message-queue/tests/rabbitmq.conf new file mode 100644 index 0000000..da58e87 --- /dev/null +++ b/homework/05-message-queue/tests/rabbitmq.conf @@ -0,0 +1,6 @@ +heartbeat = 5 +# Refresh queue state before paused workers reach the heartbeat timeout. +collect_statistics_interval = 1000 +# Allow basic solutions to use transient queues before adding fault tolerance. +# https://www.rabbitmq.com/docs/queues#durability +deprecated_features.permit.transient_nonexcl_queues = true diff --git a/homework/05-message-queue/tests/requirements.txt b/homework/05-message-queue/tests/requirements.txt new file mode 100644 index 0000000..412efa8 --- /dev/null +++ b/homework/05-message-queue/tests/requirements.txt @@ -0,0 +1,4 @@ +pytest==8.4.2 +requests==2.32.5 +docker==7.1.0 +loguru==0.7.3 diff --git a/homework/05-message-queue/tests/test_server.py b/homework/05-message-queue/tests/test_server.py new file mode 100644 index 0000000..db12cf5 --- /dev/null +++ b/homework/05-message-queue/tests/test_server.py @@ -0,0 +1,583 @@ +import io +import sys +import tarfile +import docker +import pytest +import requests +import subprocess +import time +import uuid + +from contextlib import contextmanager +from pathlib import Path +from urllib.parse import quote +from loguru import logger + + +IMAGES_ENDPOINT = 'http://localhost:5000/api/v1.0/images' +BROKER_ENDPOINT = 'http://guest:guest@localhost:15672/api' +TASK_QUEUE_ENDPOINT = f'{BROKER_ENDPOINT}/queues/%2F/task_queue' + +logger.remove() +logger.add(sys.stderr, colorize=False, format="=== TEST ===| {time:YYYY-MM-DD HH:mm:ss.SSS} {level} {message}") + + +# Tests =============================================================================================================== + +@pytest.mark.parametrize("services", [['rabbitmq', 'server']]) +def test_empty_data_dir(docker_tester): + logger.info(f"Sending GET {IMAGES_ENDPOINT}") + try: + response = requests.get(IMAGES_ENDPOINT, timeout=3) + except Exception as e: + logger.error(f"Request failed: {e}") + pytest.fail(f"Failed to get images: {e}") + logger.info(f"Got response: {response.status_code} {response.text.rstrip()}") + assert response.status_code == 200 + assert 'image_ids' in response.json() + assert len(response.json()['image_ids']) == 0 + + +@pytest.mark.parametrize("services", [['rabbitmq', 'server']]) +def test_bad_request(docker_tester): + invalid_bodies = [{}, {'image_url': None}, {'image_url': 1}, + {'image_url': ''}, [], None, 'image-url'] + for body in invalid_bodies: + logger.info(f"Sending invalid POST: {body!r}") + response = requests.post(IMAGES_ENDPOINT, json=body, timeout=3) + assert response.status_code == 400, response.text + response = requests.post(IMAGES_ENDPOINT, data='{', + headers={'Content-Type': 'application/json'}, timeout=3) + assert response.status_code == 400, response.text + + +@pytest.mark.parametrize("services", [['rabbitmq', 'server']]) +def test_nonexistent_image(docker_tester): + nonexistent_image_id = str(uuid.uuid4()) + logger.info(f"Sending GET {IMAGES_ENDPOINT}/{nonexistent_image_id}") + try: + response = requests.get(f'{IMAGES_ENDPOINT}/{nonexistent_image_id}', timeout=3) + except Exception as e: + logger.error(f"Request failed: {e}") + pytest.fail(f"Failed to get images: {e}") + logger.info(f"Got response: {response.status_code} {response.text.rstrip()}") + assert response.status_code == 404 + + +@pytest.mark.parametrize("services", [['rabbitmq', 'server']]) +def test_task_queue(docker_tester): + time.sleep(5) + check_task_queue(0, 10) + pending_ids = post_images(10) + check_unprocessed_images(pending_ids, docker_tester) + check_task_queue(10, 10) + time.sleep(5) + check_task_queue(10, 10) + check_unprocessed_images(pending_ids, docker_tester) + + +@pytest.mark.parametrize("services", [['rabbitmq', 'server', 'worker']]) +def test_single_image(docker_tester): + pending_ids = post_images(1) + wait_and_check_results(pending_ids, 10, docker_client=docker_tester) + + +@pytest.mark.parametrize("services", [['rabbitmq', 'server', 'worker']]) +def test_multiple_images(docker_tester): + pending_ids = post_images(10) + wait_and_check_results(pending_ids, 10, docker_client=docker_tester) + + +@pytest.mark.parametrize("services", [['rabbitmq', 'server', 'worker']]) +def test_captions_generated_on_workers(docker_tester): + worker1 = docker_tester.containers.get("distsys-mq-worker-1") + worker2 = docker_tester.containers.get("distsys-mq-worker-2") + wait_for_worker_ready(worker1) + wait_for_worker_ready(worker2) + worker1.pause() + worker2.pause() + + pending_ids = post_images(10) + time.sleep(5) + + check_unprocessed_images(pending_ids, docker_tester) + worker1.unpause() + worker2.unpause() + wait_and_check_results(pending_ids, 10, docker_client=docker_tester) + + +@pytest.mark.parametrize("services", [['rabbitmq', 'server-fdv', 'worker']]) +def test_multiple_images_no_listdir(docker_tester): + pending_ids = post_images(10) + # This server has no access to the workers' result volume. + wait_and_check_results(pending_ids, 10, check_captions=False, docker_client=docker_tester) + + +@pytest.mark.parametrize("services", [['rabbitmq', 'server', 'worker']]) +def test_heartbeats_timeout(docker_tester): + # With heartbeat=5, closing a connection without heartbeats can take 15s. + time.sleep(20) + pending_ids = post_images(10) + wait_and_check_results(pending_ids, 10, docker_client=docker_tester) + + +@pytest.mark.parametrize("services", [['rabbitmq', 'server', 'worker']]) +def test_publisher_confirms(docker_tester): + rabbit = docker_tester.containers.get("distsys-mq-rabbitmq-1") + rabbit.pause() + pending_ids = post_images(10) + time.sleep(5) + rabbit.kill() + time.sleep(1) + rabbit.start() + wait_for_task_queue(lambda queue: True, "broker restart", max_attempts=20) + wait_and_check_results(pending_ids, 10, docker_client=docker_tester) + + +@pytest.mark.parametrize("services", [['rabbitmq', 'server', 'worker']]) +def test_faulty_worker(docker_tester): + worker1 = docker_tester.containers.get("distsys-mq-worker-1") + worker2 = docker_tester.containers.get("distsys-mq-worker-2") + # Prevent the healthy worker from completing the whole batch before the fault. + worker2.kill() + subscribed = wait_for_worker_ready(worker1) + worker1.pause() + pending_ids = post_images(10) + wait_for_pending_tasks(subscribed) + worker1.kill() + worker2.start() + wait_and_check_results(pending_ids, 10, docker_client=docker_tester) + + +@pytest.mark.parametrize("services", [['rabbitmq', 'server', 'worker']]) +def test_two_faulty_workers(docker_tester): + worker1 = docker_tester.containers.get("distsys-mq-worker-1") + worker2 = docker_tester.containers.get("distsys-mq-worker-2") + subscribed1 = wait_for_worker_ready(worker1) + subscribed2 = wait_for_worker_ready(worker2) + worker1.pause() + worker2.pause() + pending_ids = post_images(10) + wait_for_pending_tasks(subscribed1 or subscribed2) + worker1.kill() + worker2.kill() + worker1.start() + wait_and_check_results(pending_ids, 10, docker_client=docker_tester) + + +@pytest.mark.parametrize("services", [['rabbitmq', 'server', 'worker']]) +def test_faulty_worker_and_rabbit_restart(docker_tester): + worker1 = docker_tester.containers.get("distsys-mq-worker-1") + worker2 = docker_tester.containers.get("distsys-mq-worker-2") + rabbit = docker_tester.containers.get("distsys-mq-rabbitmq-1") + # No worker may finish tasks before the broker loses its in-memory state. + worker1.kill() + worker2.kill() + pending_ids = post_images(10) + check_task_queue(10, 10) + rabbit.kill() + rabbit.start() + wait_for_task_queue(lambda queue: True, "broker restart", max_attempts=20) + worker2.start() + wait_and_check_results(pending_ids, 10, docker_client=docker_tester) + + +@pytest.mark.parametrize("services", [['rabbitmq', 'server', 'worker']]) +def test_total_eclipse_of_the_heart(docker_tester): + completed_ids = post_images(10) + wait_and_check_results(completed_ids, 10, docker_client=docker_tester) + + worker1 = docker_tester.containers.get("distsys-mq-worker-1") + worker2 = docker_tester.containers.get("distsys-mq-worker-2") + rabbit = docker_tester.containers.get("distsys-mq-rabbitmq-1") + worker1.kill() + worker2.kill() + rabbit.kill() + + # These URLs repeat the first batch, but every request needs a new ID. + pending_ids = post_images(10, known_ids=completed_ids) + # Also accept requests after the client has detected the broken connection. + time.sleep(20) + later_ids = post_images(10, known_ids=completed_ids | pending_ids) + wait_and_check_results(completed_ids, 10, docker_client=docker_tester) + + rabbit.start() + wait_for_task_queue(lambda queue: True, "broker restart", max_attempts=20) + worker1.start() + worker2.start() + wait_and_check_results(completed_ids | pending_ids | later_ids, 10, docker_client=docker_tester) + + +@pytest.mark.parametrize("services", [['rabbitmq', 'server', 'worker', 'network-fault']]) +def test_notifications_survive_rabbit_restart(docker_tester): + worker1 = docker_tester.containers.get("distsys-mq-worker-1") + worker2 = docker_tester.containers.get("distsys-mq-worker-2") + rabbit = docker_tester.containers.get("distsys-mq-rabbitmq-1") + fault = docker_tester.containers.get("distsys-mq-network-fault-1") + worker1.kill() + worker2.kill() + pending_ids = post_images(10) + check_task_queue(10, 10) + + rabbit.reload() + address = next(network['IPAddress'] for network in + rabbit.attrs['NetworkSettings']['Networks'].values() if network['IPAddress']) + with block_amqp_link(fault, address): + worker1.start() + worker2.start() + wait_for_queued_notifications() + # Notifications cannot reach the API; tasks may still await acknowledgement. + wait_and_check_results(set(), 1, docker_client=docker_tester) + worker1.kill() + worker2.kill() + rabbit.kill() + rabbit.start() + wait_for_task_queue(lambda queue: True, "broker restart", max_attempts=20) + + worker1.start() + wait_and_check_results(pending_ids, 10, docker_client=docker_tester) + + +@pytest.mark.parametrize("services", [['rabbitmq', 'server', 'worker', 'network-fault']]) +def test_duplicate_deliveries(docker_tester): + worker1 = docker_tester.containers.get("distsys-mq-worker-1") + worker2 = docker_tester.containers.get("distsys-mq-worker-2") + server = docker_tester.containers.get("distsys-mq-server-1") + rabbit = docker_tester.containers.get("distsys-mq-rabbitmq-1") + fault = docker_tester.containers.get("distsys-mq-network-fault-1") + worker1.kill() + worker2.kill() + pending_ids = post_images(10) + check_task_queue(10, 10) + task_copies = duplicate_queued_messages('task_queue') + assert task_copies == 20 + check_task_queue(task_copies, 10) + + rabbit.reload() + address = next(network['IPAddress'] for network in + rabbit.attrs['NetworkSettings']['Networks'].values() if network['IPAddress']) + with block_amqp_link(fault, address): + # Close the blocked connection so in-flight notifications become ready. + disconnect_from_broker(server) + worker1.start() + worker2.start() + queues = wait_for_queued_notifications(ready_only=True) + wait_and_check_results(set(), 1, docker_client=docker_tester) + worker1.kill() + worker2.kill() + notification_copies = { + queue['name']: duplicate_queued_messages(queue['name']) + for queue in queues if queue['name'] != 'task_queue' and queue.get('messages', 0) > 0 + } + assert notification_copies + wait_for_queue_counts(notification_copies) + + worker1.start() + wait_and_check_results(pending_ids, 10, docker_client=docker_tester) + # The first complete result list may precede processing of later copies. + wait_for_queue_counts({name: 0 for name in notification_copies}) + wait_and_check_results(pending_ids, 1, docker_client=docker_tester) + + +# Utils =============================================================================================================== + +@pytest.fixture +def docker_tester(services): + print() + try: + run_docker_compose_up(services) + check_server_endpoint() + client = docker.from_env() + yield client + finally: + print() + run_docker_compose_down() + + +def run_docker_compose_up(services): + command = ["docker", "compose", "--ansi", "never", "up", "--force-recreate"] + for service in services: + command.append(service) + subprocess.Popen(command, cwd=Path(__file__).parent.parent.absolute(), stdout=None, stderr=None) + + +def run_docker_compose_down(): + command = ["docker", "compose", "down", "--volumes"] + subprocess.run(command, cwd=Path(__file__).parent.parent.absolute(), stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL) + + +def check_server_endpoint(max_attempts=20): + attempt = 0 + while True: + attempt += 1 + try: + logger.info(f"Sending GET {IMAGES_ENDPOINT}") + response = requests.get(IMAGES_ENDPOINT, timeout=3) + assert response.status_code == 200, response.text + logger.info(f"Attempt {attempt} succeeded: got response, server endpoint is ready") + return + except Exception as e: + logger.error(f"Attempt {attempt} failed: {e}") + if attempt == max_attempts: + logger.error(f"Max attempts reached, give up") + pytest.fail("Server endpoint is not ready") + logger.info(f"Retry in 3 seconds...") + time.sleep(3) + + +def post_images(num_requests, known_ids=None): + known_ids = set() if known_ids is None else known_ids + pending_ids = set() + for i in range(num_requests): + input_data = {"image_url": f"https://somehost.com/some-image-{i}.jpg"} + logger.info(f"Sending POST {IMAGES_ENDPOINT} {input_data}") + try: + response = requests.post(IMAGES_ENDPOINT, json=input_data, timeout=3) + except Exception as e: + logger.error(f"Request failed: {e}") + pytest.fail(f"Failed to post image: {e}") + logger.info(f"Got response: {response.status_code} {response.text.rstrip()}") + assert response.status_code == 200 + assert 'image_id' in response.json() + image_id = response.json()['image_id'] + assert type(image_id) in (str, int), 'image_id must be a string or integer' + assert image_id not in known_ids, f'Reused image_id: {image_id}' + assert image_id not in pending_ids + pending_ids.add(image_id) + return pending_ids + + +def wait_and_check_results(pending_ids, max_attempts, docker_client, check_captions=True): + expected_count = len(pending_ids) + attempt = 0 + while True: + attempt += 1 + try: + logger.info(f"Sending GET {IMAGES_ENDPOINT}") + response = requests.get(IMAGES_ENDPOINT, timeout=3) + logger.info(f"Got response: {response.status_code} {response.text.rstrip()}") + assert response.status_code == 200 + assert 'image_ids' in response.json() + ready_ids = response.json()['image_ids'] + assert isinstance(ready_ids, list), 'image_ids must be a list' + assert all(type(image_id) in (str, int) for image_id in ready_ids) + ready_set = set(ready_ids) + assert len(ready_ids) == len(ready_set), 'image_ids contains duplicates' + assert ready_set <= pending_ids, 'image_ids contains unexpected IDs' + count = len(ready_set) + if ready_set == pending_ids: + if check_captions: + for image_id in pending_ids: + check_image_caption(image_id, docker_client) + logger.info(f"Attempt {attempt} succeeded: got {count} results as expected") + return + else: + logger.info(f"Attempt {attempt} not succeeded: got {count} results, expect {expected_count}") + except Exception as e: + logger.error(f"Attempt {attempt} failed: {e}") + if attempt == max_attempts: + logger.error(f"Max attempts reached, give up") + pytest.fail("Server didn't return expected results") + logger.info(f"Retry in 3 seconds...") + time.sleep(3) + + +def check_image_caption(image_id, docker_client): + logger.info(f"Sending GET {IMAGES_ENDPOINT}/{image_id}") + try: + response = requests.get(f'{IMAGES_ENDPOINT}/{image_id}', timeout=3) + except Exception as e: + logger.error(f"Request failed: {e}") + pytest.fail(f"Failed to get image caption: {e}") + logger.info(f"Got response: {response.status_code} {response.text.rstrip()}") + assert response.status_code == 200 + assert 'caption' in response.json() + assert isinstance(response.json()['caption'], str) + saved_caption = read_saved_caption(image_id, docker_client) + assert response.json()['caption'] == saved_caption, 'API caption differs from the saved result' + + +def read_saved_caption(image_id, docker_client): + # Docker's archive API works even if the solution image has no shell or Python. + server = docker_client.containers.get('distsys-mq-server-1') + stream, _ = server.get_archive(f'/data/{image_id}.txt') + with tarfile.open(fileobj=io.BytesIO(b''.join(stream))) as archive: + files = [member for member in archive.getmembers() if member.isfile()] + assert len(files) == 1, 'Expected one regular caption file' + with archive.extractfile(files[0]) as result: + return result.read().decode('utf-8') + + +def check_unprocessed_images(image_ids, docker_client): + wait_and_check_results(set(), 1, docker_client) + for image_id in image_ids: + response = requests.get(f'{IMAGES_ENDPOINT}/{image_id}', timeout=3) + assert response.status_code == 404, f'Result available without workers: {response.text}' + + +def wait_for_task_queue(predicate, description, max_attempts=10, interval=3): + last_state = None + for attempt in range(1, max_attempts + 1): + try: + response = requests.get(TASK_QUEUE_ENDPOINT, timeout=3) + assert response.status_code == 200, response.text + queue = response.json() + last_state = {key: queue.get(key) for key in + ('messages_ready', 'messages_unacknowledged', 'consumers', 'message_stats')} + if predicate(queue): + logger.info(f"Task queue is ready: {description}") + return queue + except Exception as e: + last_state = str(e) + logger.info(f"Waiting for {description}: {last_state}") + if attempt < max_attempts: + time.sleep(interval) + pytest.fail(f"Task queue did not reach {description}; last state: {last_state}") + + +def check_task_queue(expected_count, max_attempts): + return wait_for_task_queue(lambda queue: queue['messages_ready'] == expected_count, + f'{expected_count} ready messages', max_attempts) + + +def wait_for_worker_ready(worker, max_attempts=10): + last_state = None + for attempt in range(1, max_attempts + 1): + try: + worker.reload() + addresses = {network['IPAddress'] for network in + worker.attrs['NetworkSettings']['Networks'].values() if network['IPAddress']} + response = requests.get(TASK_QUEUE_ENDPOINT, timeout=3) + assert response.status_code == 200, response.text + consumers = response.json().get('consumer_details', []) + if any(consumer['channel_details']['peer_host'] in addresses for consumer in consumers): + return True + # Polling with basic.get is also valid; it has no registered consumer. + response = requests.get(f'{BROKER_ENDPOINT}/channels', timeout=3) + assert response.status_code == 200, response.text + channels = response.json() + last_state = [{'peer_host': channel['connection_details']['peer_host'], + 'message_stats': channel.get('message_stats')} for channel in channels + if channel['connection_details']['peer_host'] in addresses] + for channel in channels: + if channel['connection_details']['peer_host'] in addresses: + stats = channel.get('message_stats') or {} + if any(stats.get(name, 0) > 0 for name in ('get', 'get_no_ack', 'get_empty')): + return False + except Exception as e: + last_state = str(e) + if attempt < max_attempts: + time.sleep(3) + pytest.fail(f"Worker {worker.name} did not start consuming tasks; last state: {last_state}") + + +def wait_for_pending_tasks(subscribed): + if subscribed: + # Delivery statistics can arrive before the unacknowledged count. + # Automatic ACK is deliberately allowed through this barrier: after the + # kill, the missing results must expose the lost deliveries instead. + predicate = lambda queue: (queue['messages_unacknowledged'] > 0 or + any((queue.get('message_stats') or {}).get(name, 0) > 0 + for name in ('deliver', 'deliver_no_ack'))) + description = 'a task delivered to a paused worker' + else: + predicate = lambda queue: queue['messages_ready'] == 10 + description = '10 tasks waiting for a polling worker' + # Statistics refresh every 1s in rabbitmq.conf. + return wait_for_task_queue(predicate, description, max_attempts=10, interval=1) + + +def wait_for_queued_notifications(max_attempts=20, ready_only=False): + last_state = None + for attempt in range(1, max_attempts + 1): + try: + response = requests.get(f'{BROKER_ENDPOINT}/queues/%2F', timeout=3) + assert response.status_code == 200, response.text + queues = response.json() + last_state = {queue['name']: queue.get('messages') for queue in queues} + # Notifications may be batched or spread across multiple queues. + # Workers may keep tasks unacknowledged until the API receives a result. + notification_count = sum(count or 0 for name, count in last_state.items() + if name != 'task_queue') + ready = sum(queue.get('messages_ready', 0) or 0 for queue in queues + if queue['name'] != 'task_queue') + if 'task_queue' in last_state and notification_count > 0 and (not ready_only or ready == notification_count): + return queues + except Exception as e: + last_state = str(e) + if attempt < max_attempts: + time.sleep(1) + pytest.fail(f'No queued notifications reached the required state; last state: {last_state}') + + +@contextmanager +def block_amqp_link(fault, address): + rules = [('INPUT', '-s', address, '--sport'), ('OUTPUT', '-d', address, '--dport')] + installed = [] + try: + for chain, direction, host, port in rules: + rule = [chain, direction, host, '-p', 'tcp', port, '5672', '-j', 'DROP'] + result = fault.exec_run(['iptables', '-I', *rule]) + assert result.exit_code == 0, result.output.decode(errors='replace') + installed.append(rule) + yield + finally: + for rule in reversed(installed): + result = fault.exec_run(['iptables', '-D', *rule]) + assert result.exit_code == 0, result.output.decode(errors='replace') + + +def duplicate_queued_messages(queue_name): + # Consumers/producers of this queue are stopped or disconnected by the test. + messages = [] + while True: + response = requests.post( + f'{BROKER_ENDPOINT}/queues/%2F/{quote(queue_name, safe="")}/get', + json={'count': 100, 'ackmode': 'ack_requeue_false', 'encoding': 'base64'}, + timeout=3) + assert response.status_code == 200, response.text + batch = response.json() + if not batch: + break + messages.extend(batch) + assert messages, f'No messages to duplicate in {queue_name}' + for message in messages: + for _ in range(2): + response = requests.post(f'{BROKER_ENDPOINT}/exchanges/%2F/amq.default/publish', + json={'routing_key': queue_name, 'properties': message['properties'], + 'payload': message['payload'], 'payload_encoding': 'base64'}, timeout=3) + assert response.status_code == 200, response.text + assert response.json()['routed'] is True, f'Copy was not routed to {queue_name}' + logger.info(f'Queued {len(messages) * 2} copies in {queue_name}') + return len(messages) * 2 + + +def disconnect_from_broker(container): + container.reload() + addresses = {network['IPAddress'] for network in + container.attrs['NetworkSettings']['Networks'].values() if network['IPAddress']} + response = requests.get(f'{BROKER_ENDPOINT}/connections', timeout=3) + assert response.status_code == 200, response.text + for connection in response.json(): + if connection['peer_host'] in addresses: + response = requests.delete( + f'{BROKER_ENDPOINT}/connections/{quote(connection["name"], safe="")}', timeout=3) + assert response.status_code in (204, 404), response.text + + +def wait_for_queue_counts(expected, max_attempts=30): + last_state = None + for attempt in range(1, max_attempts + 1): + try: + last_state = {} + for name in expected: + response = requests.get(f'{BROKER_ENDPOINT}/queues/%2F/{quote(name, safe="")}', timeout=3) + assert response.status_code == 200, response.text + last_state[name] = response.json().get('messages') + if last_state == expected: + return + except Exception as e: + last_state = str(e) + if attempt < max_attempts: + time.sleep(1) + pytest.fail(f'Queues did not reach {expected}; last state: {last_state}')