Metadata-Version: 2.1
Name: wexample-queue
Version: 1.1.0
Summary: Manage Queuing
Author-Email: weeger <contact@wexample.com>
License: MIT
Classifier: Programming Language :: Python :: 3
Classifier: License :: OSI Approved :: MIT License
Classifier: Operating System :: OS Independent
Requires-Python: >=3.10
Requires-Dist: pika>=1.3.0
Requires-Dist: wexample-app>=19.2.0
Requires-Dist: wexample-helpers>=20.0.0
Provides-Extra: dev
Requires-Dist: pytest; extra == "dev"
Requires-Dist: pytest-cov; extra == "dev"
Description-Content-Type: text/markdown

# queue

Version: 1.1.0

`wexample-queue` is what a worker needs to live on a RabbitMQ queue and nothing more: a connector holding the broker connection, and a service base whose whole subject is one queue. A subclass writes the name of the queue and what to do with a message; connecting, declaring, decoding, acknowledging and giving up are the same for every worker and are already written.

It exists because the other end of the queue is not always Python. The body of a message is plain JSON, published on the default exchange with the queue name as routing key — which is where a Symfony Messenger transport declaring that same queue binds its own exchange, so both ends meet in the same place without either knowing the other.

## A worker

```python
from wexample_queue.service.abstract_queue_service import AbstractQueueService

@base_class
class ProcessRunService(AbstractQueueService):
    queue_name: str = public_field(
        default="process_run",
        description="Name of the queue this service consumes",
    )

    def process_message(self, message: dict[str, Any]) -> None:
        self.log(f"Run {message['id']} asked for in {message['workdir']}")
```

Calling `start()` opens the connection, declares the queue and blocks consuming
until `stop()`. Returning from `process_message()` acknowledges the message;
raising rejects it without putting it back, because requeueing a message that
just failed is asking for the same failure at once and for ever. Where it goes
then is the broker's business — a dead letter exchange if the queue has one,
nowhere if it has not.

## Publishing

```python
self.connector.publish("process_run_event", {"kind": "process_run", "id": run_id})
```

## The broker

Where the broker is comes from the environment, under four keys the connector
declares as expected:

```
RABBITMQ_HOST
RABBITMQ_PORT
RABBITMQ_USER
RABBITMQ_PASSWORD
```

Reading them from there rather than from a caller is what keeps a worker and
whoever publishes to it from holding two addresses that could disagree.

## Table of Contents

- [A worker](#a-worker)
- [Publishing](#publishing)
- [The broker](#the-broker)
- [Installation](#installation)
- [Tests](#tests)
- [Architecture](#architecture)
- [Integration in the Suite](#integration-in-the-suite)
- [Dependencies](#dependencies)
- [Versioning & Compatibility Policy](#versioning--compatibility-policy)
- [License](#license)
- [About us](#about-us)
- [Known Limitations & Roadmap](#known-limitations--roadmap)
- [Status & Compatibility](#status--compatibility)
- [Useful Links](#useful-links)
- [Migration Notes](#migration-notes)

## Installation

```bash
pip install wexample-queue
```

Requires Python >=3.10.

```bash
pip install wexample-queue
```

A worker is a class with a queue name and a `process_message()`:

```python
from typing import Any

from wexample_helpers.classes.field import public_field
from wexample_helpers.decorator.base_class import base_class
from wexample_queue.service.abstract_queue_service import AbstractQueueService

@base_class
class HelloService(AbstractQueueService):
    queue_name: str = public_field(
        default="hello",
        description="Name of the queue this service consumes",
    )

    def process_message(self, message: dict[str, Any]) -> None:
        self.log(f"Got {message}")

HelloService(kernel=kernel, queue_name="hello").start()
```

`start()` blocks. The four `RABBITMQ_*` environment keys have to be readable by
the kernel, which is where the connector goes to find the broker.

## Tests

This project uses `pytest` for testing and `pytest-cov` for code coverage analysis.

### Installation

First, install the required testing dependencies:
```bash
.venv/bin/python -m pip install pytest pytest-cov
```

### Basic Usage

Run all tests with coverage:
```bash
.venv/bin/python -m pytest --cov --cov-report=html
```

### Common Commands
```bash
# Run tests with coverage for a specific module
.venv/bin/python -m pytest --cov=your_module

# Show which lines are not covered
.venv/bin/python -m pytest --cov=your_module --cov-report=term-missing

# Generate an HTML coverage report
.venv/bin/python -m pytest --cov=your_module --cov-report=html

# Combine terminal and HTML reports
.venv/bin/python -m pytest --cov=your_module --cov-report=term-missing --cov-report=html

# Run specific test file with coverage
.venv/bin/python -m pytest tests/test_file.py --cov=your_module --cov-report=term-missing
```

### Viewing HTML Reports

After generating an HTML report, open `htmlcov/index.html` in your browser to view detailed line-by-line coverage information.

### Coverage Threshold

To enforce a minimum coverage percentage:
```bash
.venv/bin/python -m pytest --cov=your_module --cov-fail-under=80
```

This will cause the test suite to fail if coverage drops below 80%.

## Architecture

Two classes, and the line between them is what a broker is against what a worker is.

### The connector

src/wexample_queue/connector/rabbitmq_external_connector.py extends `AbstractExternalConnector` from `wexample-app` and holds a `pika.BlockingConnection` and its channel. It exposes only what this package needs — `connect`, `disconnect`, `declare_queue`, `publish`, `consume`, `stop_consuming` — and reads the broker address from four environment keys it declares in `get_expected_env_keys()`, so a worker never carries one.

`publish()` goes through the default exchange with the queue name as routing key, and marks the message persistent. `declare_queue()` declares it durable, with the same arguments both ends use: neither side depends on the other having started first, and a broker restarted on its own finds its queues again.

### The service

src/wexample_queue/service/abstract_queue_service.py extends `AbstractService` from `wexample-app`, which already brings the run loop, the signal handling and the shutdown. What this adds is the queue: `start()` opens the connection and declares the queue before handing over to the loop, `_run()` blocks consuming, and `_on_message()` decodes one delivery and answers the broker for it.

The answer is the part worth reading. A body that is not JSON, and a message `process_message()` raised on, are both rejected with `requeue=False`. Putting such a message back would be asking for the same failure immediately and for ever; dropping it silently with an `ack` would lose it without saying so. Rejecting it hands the decision to the broker, which is the only place a dead letter exchange can be configured.

Unlike the implementation this was drawn from, constructing a service does not start it. `start()` is called by whoever owns the worker, which is what makes a service testable without a broker.

## Integration in the Suite

This package is part of the Wexample Suite — a collection of high-quality, modular tools designed to work seamlessly together across multiple languages and environments.

### Related Packages

The suite includes packages for configuration management, file handling, prompts, and more. Each package can be used independently or as part of the integrated suite.

Visit the [Wexample Suite documentation](https://docs.wexample.com) for the complete package ecosystem.

## Dependencies

- pika: >=1.3.0
- wexample-app: >=19.2.0
- wexample-helpers: >=20.0.0

## Versioning & Compatibility Policy

Wexample packages follow **Semantic Versioning** (SemVer):

- **MAJOR**: Breaking changes
- **MINOR**: New features, backward compatible
- **PATCH**: Bug fixes, backward compatible

We maintain backward compatibility within major versions and provide clear migration guides for breaking changes.

## License

This project is licensed under the MIT License - see the [LICENSE](LICENSE) file for details.

Free to use in both personal and commercial projects.

## About us

[Wexample](https://wexample.com) stands as a cornerstone of the digital ecosystem — a collective of seasoned engineers, researchers, and creators driven by a relentless pursuit of technological excellence. More than a media platform, it has grown into a vibrant community where innovation meets craftsmanship, and where every line of code reflects a commitment to clarity, durability, and shared intelligence.

This packages suite embodies this spirit. Trusted by professionals and enthusiasts alike, it delivers a consistent, high-quality foundation for modern development — open, elegant, and battle-tested. Its reputation is built on years of collaboration, refinement, and rigorous attention to detail, making it a natural choice for those who demand both robustness and beauty in their tools.

Wexample cultivates a culture of mastery. Each package, each contribution carries the mark of a community that values precision, ethics, and innovation — a community proud to shape the future of digital craftsmanship.

## Known Limitations & Roadmap

Current limitations and planned features are tracked in the GitHub issues.

See the [project roadmap](https://github.com/wexample/python-queue/issues) for upcoming features and improvements.

## Status & Compatibility

**Maturity**: Production-ready

**Python Support**: >=3.10

**OS Support**: Linux, macOS, Windows

**Status**: Actively maintained

## Useful Links

- **Homepage**: https://github.com/wexample/python-queue
- **Documentation**: [docs.wexample.com](https://docs.wexample.com)
- **Issue Tracker**: https://github.com/wexample/python-queue/issues
- **Discussions**: https://github.com/wexample/python-queue/discussions
- **PyPI**: [pypi.org/project/wexample-queue](https://pypi.org/project/wexample-queue/)

## Migration Notes

When upgrading between major versions, refer to the migration guides in the documentation.

Breaking changes are clearly documented with upgrade paths and examples.
