jiengup opened a new issue, #4019:
URL: https://github.com/apache/iggy/issues/4019
### Description
### Goal
Define the complete Python producer API and lifecycle model, then deliver a
fully functional direct producer. The public API introduced here should remain
stable when background mode is added. **To make this issue clear, Background
Mode will be solved in a separate issue.**
The Python SDK's behavior must follow the Rust core. A documentation
reference: #3989. Will review this document PR when implementing this issue.
### TODO
- Add a new producer module containing:
- `IggyProducer`
- `DirectProducerConfig`
- `BackgroundProducerConfig`
- `ProducerSharding`
- `BackpressureMode`
- Add async `IggyClient.producer()`.
- Build and initialize the Rust producer before returning it.
- Support:
- Default partitioning
- Stream creation if missing
- Topic creation if missing
- Topic partitions, expiry, and maximum size
- Send retry count and interval
- Expose:
- `send()`
- `send_one()`
- `send_with_partitioning()`
- `send_to()`
- `shutdown()`
- Add async context-manager support.
- Preserve the existing low-level `IggyClient.send_messages()` API.
- Export all new types from the Python module.
- Update `apache_iggy.pyi`.
- Add direct-producer documentation and examples.
- Add unit and integration tests for all functionality delivered in this PR.
### Proposed API
```python
producer = await client.producer(
stream,
topic,
partitioning=None,
mode=DirectProducerConfig(),
create_stream_if_not_exists=True,
create_topic_if_not_exists=True,
topic_partitions_count=1,
topic_message_expiry=None,
topic_max_size=None,
send_retries=3,
send_retry_interval=timedelta(seconds=1),
)
```
```python
async with await client.producer("orders", "created") as producer:
await producer.send(messages)
```
### Lifecycle design
Store the Rust producer as:
```rust
Arc<RwLock<Option<RustIggyProducer>>>
```
This allows concurrent sends while preserving the ownership required by
Rust’s consuming `shutdown(self)` method.
The lifecycle semantics should be:
- `IggyClient.producer()` calls `build()` and `init()`.
- A returned producer is immediately ready to send.
- `shutdown()` waits for active sends and takes ownership of the Rust
producer.
- Repeated `shutdown()` calls are idempotent.
- Sending after shutdown raises a clear `RuntimeError`.
- `__aexit__()` calls `shutdown()`.
- Async cleanup must not depend on `__del__`.
### Acceptance criteria
- A direct producer can be created and used from Python.
- Its defaults match the Rust builder defaults.
- All four send methods return `SendMessagesResponse`.
- Default and per-call partitioning work correctly.
- Missing streams and topics are created according to the configured flags.
- Topic creation options are propagated correctly.
- Disabling automatic creation produces an error when the resource is
missing.
- Default, custom, and disabled retry policies work as expected.
- Direct batching and linger configuration are propagated to Rust.
- Empty-message behavior matches the Rust SDK.
- Shutdown is idempotent, and sending afterward fails clearly.
- The async context manager shuts the producer down on normal and
exceptional exit.
- Initialization failure does not return a partially initialized producer.
- Invalid types, durations, and numeric ranges produce stable Python
exceptions.
- `.pyi`, README documentation, and a complete direct-producer example are
included.
- Existing `IggyClient.send_messages()` behavior remains unchanged.
- Formatting, linting, stub-generation checks, Rust tests, and Python tests
pass.
---
## PR 2: Add Background Producer, Backpressure, and Sharding
### Goal
Implement background sending on top of the API and lifecycle framework
established by PR 1.
### TODO
- Convert `BackgroundProducerConfig` into the Rust `BackgroundConfig`.
- Support:
- `num_shards`
- `linger_time`
- `batch_size`
- `batch_length`
- `max_buffer_size`
- `max_in_flight`
- Implement:
- `ProducerSharding.ORDERED`
- `ProducerSharding.BALANCED`
- Implement:
- `BackpressureMode.block()`
- `BackpressureMode.block_with_timeout()`
- `BackpressureMode.fail_immediately()`
- Connect `IggyProducer.send*()` to the Rust background dispatcher.
- Ensure `shutdown()` flushes accepted buffered messages.
- Complete background-specific `.pyi` documentation and examples.
- Add unit and integration tests for background behavior.
A custom Python sharding callback and Python error callback are not part of
this PR. Ordered and balanced Rust implementations are sufficient.
### Acceptance criteria
- A producer can be created with `BackgroundProducerConfig`.
- All background configuration fields are propagated correctly.
- `send()`, `send_one()`, `send_with_partitioning()`, and `send_to()` work
in background mode.
- A successful background send means the dispatcher accepted the messages,
not that the server committed them.
- Background sends return an empty confirmation list, matching the Rust SDK.
- Batch length, batch size, and linger-time flush conditions are covered by
tests.
- `ORDERED` sharding preserves ordering for the same destination.
- `BALANCED` sharding distributes work across shards and documents that
ordering is not guaranteed.
- All three backpressure modes are tested, including timeout behavior.
- Buffer and in-flight limits are enforced.
- `shutdown()` flushes all accepted messages before returning.
- Dropping a producer without shutdown is documented as potentially losing
buffered messages.
- Repeated shutdown remains safe.
- Concurrent send and shutdown behavior is tested and does not deadlock.
- The async context manager flushes background messages on normal and
exceptional exit.
- Retry and reconnect behavior is covered by integration tests.
- `.pyi`, README documentation, and a complete background-producer example
are included.
- Formatting, linting, stub-generation checks, Rust tests, and Python tests
pass.
### Contribution
- [ ] I'm willing to submit a pull request to implement this feature
### Good first issue
- [ ] I think this could be a good first issue for a new contributor
### Affected area / component
_No response_
### Proposed solution
_No response_
### Alternatives considered
_No response_
### Contribution
- [ ] I'm willing to submit a pull request to implement this feature
### Good first issue
- [ ] I think this could be a good first issue for a new contributor
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]