Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -141,7 +141,7 @@ Per-subscriber knobs (passed to `@broker.subscriber("topic", ...)`):
re-deliver the timer (duplicate). Increase if your handlers are slow.
- `polling_interval` (default `0.05` s) — base poll interval used when the topic
has due timers or just transitioned from idle. Doubles on each consecutive empty
cycle, capped at `max_polling_interval`, with ±50% jitter applied each sleep.
cycle, capped at `max_polling_interval`, with ±50% jitter applied before the cap.
- `max_polling_interval` (default `5.0` s) — ceiling for the adaptive idle
backoff. Lower it for tighter delivery latency on idle topics; raise it to
reduce Redis load on workloads with long idle stretches.
Expand Down
4 changes: 2 additions & 2 deletions docs/introduction/how-it-works.md
Original file line number Diff line number Diff line change
Expand Up @@ -64,12 +64,12 @@ The default ack policy is `NACK_ON_ERROR`: the timer is acknowledged on success,
| Parameter | Default | Description |
|-----------|---------|-------------|
| `polling_interval` | `0.05` s | Base poll interval used when the topic has due timers or just transitioned from idle |
| `max_polling_interval` | `5.0` s | Ceiling for the adaptive idle backoff — `polling_interval` doubles up to this value on consecutive empty polls. Worst-case delivery latency on a previously-idle topic is `max_polling_interval × 1.5` (with ±50% jitter) |
| `max_polling_interval` | `5.0` s | Ceiling for the adaptive idle backoff — `polling_interval` doubles up to this value on consecutive empty polls. Jitter is applied before the cap, so worst-case delivery latency on a previously-idle topic is `max_polling_interval` |
| `max_concurrent` | `5` | Max handlers running concurrently per subscriber; also bounds fetch batch size |
| `lease_ttl` | `30` s | How long a worker holds the lease before another worker may re-claim. Set to ~3–5× the P99 handler runtime: lower values speed up recovery from worker death, higher values tolerate handler GC pauses and clock skew |

## Operational requirements

`faststream-redis-timers` is designed for a **single-primary Redis** (Sentinel-managed primary/replica setups are supported; Redis Cluster is not — see the broker constructor docstring).
`faststream-redis-timers` is designed for a **single-primary Redis** (Sentinel-managed primary/replica setups are supported; Redis Cluster is not: constructing `TimersBroker` with a `RedisCluster` client raises `TypeError`).

Multiple brokers polling the same Redis derive timer due-times and lease deadlines from each broker's local wall clock. Keep all broker hosts NTP-synchronised: clock skew larger than `lease_ttl` between brokers can cause a clock-fast broker to re-lease a timer that another broker is still actively processing, producing duplicate delivery to handlers. The default `lease_ttl=30s` tolerates seconds of NTP drift; tune `lease_ttl` upward if your environment cannot guarantee sub-second sync.
2 changes: 1 addition & 1 deletion docs/introduction/installation.md
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@

## Requirements

- Python 3.13+
- Python 3.11+
- Redis 5.0+
- A running Redis instance

Expand Down
8 changes: 4 additions & 4 deletions docs/usage/basic.md
Original file line number Diff line number Diff line change
Expand Up @@ -81,8 +81,8 @@ async def schedule_reminder() -> None:
| Parameter | Default | Description |
|-----------|---------|-------------|
| `client` | _required for production_ | `redis.asyncio.Redis` client instance — the caller owns its lifecycle (the broker does not close it). The constructor accepts `None`, but a broker without a client raises `IncorrectState` on the first operation; `None` is the test-broker shape (see [Testing](./testing.md)). |
| `timeline_key` | `timers_timeline` | Sorted set key name |
| `payloads_key` | `timers_payloads` | Hash key name |
| `timeline_key` | `timers_timeline` | Sorted set key prefix; each topic uses `{timeline_key}:{topic}` |
| `payloads_key` | `timers_payloads` | Hash key prefix; each topic uses `{payloads_key}:{topic}` |
| `start_timeout` | `3.0` | Seconds to wait for the subscriber's first Redis ping during startup |
| `graceful_timeout` | `15.0` | Seconds to wait for in-flight timers on shutdown |

Expand All @@ -92,7 +92,7 @@ connection is released when the application stops.

## Timer IDs

Each timer has a unique `timer_id`. If you don't provide one, a UUID is generated automatically. You can supply your own to make a timer idempotent — publishing the same `timer_id` twice will **overwrite the first** silently (no error, no warning). This is the right behavior for idempotent retry of `publish()` calls but a footgun if two unrelated callers pick the same ID. Namespace your IDs (e.g., `f"invoice-{invoice_id}-due"`) to avoid accidental collisions.
Each timer has a unique `timer_id`. If you don't provide one, a UUID is generated automatically. You can supply your own to make a timer idempotent — publishing the same `timer_id` twice will **overwrite the first** silently (no error, no warning). This makes retries of `publish()` safe, but two unrelated callers that pick the same ID will overwrite each other. Namespace your IDs (e.g., `f"invoice-{invoice_id}-due"`) to avoid accidental collisions.

```python
await broker.publish(
Expand Down Expand Up @@ -166,7 +166,7 @@ If you are polling `has_pending` to detect *"the timer fired and the handler fin

### `cancel_all` race with executing handlers

If `cancel_all(topic)` runs while a worker is mid-handler for a leased timer on that topic, the handler runs to completion. When it finishes, its commit (the `ZREM` + `HDEL` that normally removes the timer) becomes a no-op because `cancel_all` has already deleted both keys for the topic. The work is *not* rolled back — only the bookkeeping is skipped — so handlers that have side effects (sent emails, written rows) will have already done them. Use `cancel_all` for topic resets, not for "stop everything in flight."
If `cancel_all(topic)` runs while a worker is mid-handler for a leased timer on that topic, the handler runs to completion. When it finishes, its commit (the `ZREM` + `HDEL` that normally removes the timer) becomes a no-op because `cancel_all` has already deleted both keys for the topic. Side effects stay; only the Redis cleanup is skipped, so handlers that have side effects (sent emails, written rows) will have already done them. Use `cancel_all` for topic resets, not for "stop everything in flight."

## Debug logging

Expand Down
1 change: 0 additions & 1 deletion docs/usage/publisher.md
Original file line number Diff line number Diff line change
Expand Up @@ -68,7 +68,6 @@ The `publish()` method on a publisher accepts the parameters below and returns t

```python
from faststream import Context
from faststream.message import StreamMessage

pub = broker.publisher("orders")

Expand Down
3 changes: 2 additions & 1 deletion docs/usage/router.md
Original file line number Diff line number Diff line change
Expand Up @@ -22,8 +22,9 @@ The `prefix` you pass to `TimersRouter` is concatenated to every topic registere
```python
router = TimersRouter(prefix="my-service:")


@router.subscriber("invoices")
async def handle_invoice(...): ...
async def handle_invoice(invoice_id: str) -> None: ...
```

…the subscriber listens on the full topic `my-service:invoices`, and Redis stores its timers under:
Expand Down
6 changes: 3 additions & 3 deletions docs/usage/subscriber.md
Original file line number Diff line number Diff line change
Expand Up @@ -66,7 +66,7 @@ Configure polling behaviour per subscriber:
| Parameter | Default | Description |
|-----------|---------|-------------|
| `polling_interval` | `0.05` s | Base poll interval; the floor used when the topic has due timers |
| `max_polling_interval` | `5.0` s | Cap for adaptive idle backoff (doubles per empty cycle, ±50% jitter) |
| `max_polling_interval` | `5.0` s | Cap for adaptive idle backoff (doubles per empty cycle, ±50% jitter applied before the cap) |
| `max_concurrent` | `5` | Max handlers running in parallel; also caps fetch batch size per poll |
| `lease_ttl` | `30` s | How long a worker holds the lease before another worker may re-claim |

Expand All @@ -81,7 +81,7 @@ Configure polling behaviour per subscriber:
async def handle_urgent(body: str) -> None: ...
```

The poll loop uses adaptive backoff: when there are no due timers, the next sleep doubles from `polling_interval` up to `max_polling_interval` and is multiplied by a random factor in `[0.5, 1.5]` to avoid thundering-herd bursts across worker fleets. The counter resets the moment a poll returns work. Worst-case delivery latency for a newly-published timer in a previously-idle topic is `max_polling_interval × 1.5`, plus any time spent waiting for the `max_concurrent` limiter when in-flight handlers are still holding capacity (back-pressure).
The poll loop uses adaptive backoff: when there are no due timers, the next sleep doubles from `polling_interval` up to `max_polling_interval` and is multiplied by a random factor in `[0.5, 1.5]` before the cap is applied, to avoid thundering-herd bursts across worker fleets. The counter resets the moment a poll returns work. Worst-case delivery latency for a newly-published timer in a previously-idle topic is `max_polling_interval`, plus any time spent waiting for the `max_concurrent` limiter when in-flight handlers are still holding capacity (back-pressure).

!!! warning "Handlers must be idempotent and concurrency-safe"
A handler that runs longer than `lease_ttl`, or a worker that crashes after the handler ran but before the commit landed, may cause the timer to be delivered more than once. Design handlers to be safe under retry. Because `max_concurrent` invocations run in parallel, handlers must also be safe under concurrent execution (no unsynchronized shared state).
Expand All @@ -97,7 +97,7 @@ The default ack policy is `NACK_ON_ERROR`: the timer is acknowledged (removed fr
| Handler returns normally | Timer removed from Redis |
| Handler raises an exception | Timer left in Redis for retry on next poll |

To manually control acknowledgement, inject the `NoCast`-typed message:
To manually control acknowledgement, inject the message:

```python
from faststream.message import StreamMessage
Expand Down
2 changes: 1 addition & 1 deletion faststream_redis_timers/configs.py
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@ def __init__(self, client: "RedisClient | None" = None) -> None:
@property
def client(self) -> "RedisClient":
if self._client is None:
msg = "Connection not available. Connect the broker first."
msg = "TimersBroker has no Redis client; pass one to TimersBroker(client)."
raise IncorrectState(msg)
return self._client

Expand Down
2 changes: 1 addition & 1 deletion tests/test_unit.py
Original file line number Diff line number Diff line change
Expand Up @@ -49,7 +49,7 @@ def test_timers_broker_registered_in_test_broker_registry() -> None:

def test_connection_state_client_not_set_raises() -> None:
state = ConnectionState()
with pytest.raises(IncorrectState):
with pytest.raises(IncorrectState, match=r"pass one to TimersBroker\(client\)"):
_ = state.client


Expand Down
Loading