Skip to content
Open
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 examples/python/README.md
Comment thread
jiengup marked this conversation as resolved.
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ To run any example, first start the server with
docker run --rm -p 8080:8080 -p 3000:3000 -p 8090:8090 apache/iggy:latest

# Or build from source (recommended for development)
cd ../../ && cargo run --bin iggy-server
cd ../../ && cargo run --bin iggy-server -- --with-default-root-credentials
```

For server configuration options and help:
Expand Down
8 changes: 6 additions & 2 deletions examples/python/basic/producer.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@
import asyncio
from typing import NamedTuple

from apache_iggy import IggyClient, StreamDetails, TopicDetails
from apache_iggy import IggyClient, Partitioning, StreamDetails, TopicDetails
from apache_iggy import SendMessage as Message
from loguru import logger

Expand Down Expand Up @@ -106,7 +106,11 @@ async def produce_messages(client: IggyClient):
await client.send_messages(
stream=STREAM_NAME,
topic=TOPIC_NAME,
partitioning=PARTITION_ID,
# A fixed strategy sends the whole batch to this partition.
# For topics with multiple partitions, Partitioning.balanced()
# distributes batches round-robin, while
# Partitioning.messages_key(key) keeps equal keys on the same partition.
partitioning=Partitioning.partition_id(PARTITION_ID),
messages=messages,
)
n_sent_batches += 1
Expand Down
111 changes: 80 additions & 31 deletions foreign/python/apache_iggy.pyi
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,7 @@ __all__ = [
"MaxTopicSize",
"OptionSpec",
"Partition",
"Partitioning",
"Permissions",
"PollingStrategy",
"ReceiveMessage",
Expand Down Expand Up @@ -872,6 +873,27 @@ class IggyClient:
Sends a ping request to the server to check connectivity.
Raises `RuntimeError` if the connection fails.
"""
def describe_options(
self, scope: builtins.str
) -> collections.abc.Awaitable[list[OptionSpec]]:
r"""
Describe the option catalog for a resource scope.

This is the discovery surface for the `options` argument on
`create_topic`/`update_topic`: a key outside the catalog is refused at
create, and the binary transports carry only the error code back.

Args:
scope: One of `"topic"`, `"stream"`, `"user"`.

Returns:
An awaitable that resolves to `list[OptionSpec]`, empty for a scope
with no keys yet.

Raises:
ValueError: If the scope name is not one of the three above.
RuntimeError: If the request fails.
"""
def login_user(
self, username: builtins.str, password: builtins.str
) -> collections.abc.Awaitable[None]:
Expand Down Expand Up @@ -1034,27 +1056,6 @@ class IggyClient:
Returns the stream details, or `None` if the stream does not exist.
Raises `RuntimeError` on failure.
"""
def describe_options(
self, scope: builtins.str
) -> collections.abc.Awaitable[builtins.list[OptionSpec]]:
r"""
Describe the option catalog for a resource scope.

This is the discovery surface for the `options` argument on
`create_topic`/`update_topic`: a key outside the catalog is refused at
create, and the binary transports carry only the error code back.

Args:
scope: One of `"topic"`, `"stream"`, `"user"`.

Returns:
An awaitable that resolves to `list[OptionSpec]`, empty for a scope
with no keys yet.

Raises:
ValueError: If the scope name is not one of the three above.
RuntimeError: If the request fails.
"""
def create_topic(
self,
stream: builtins.str | builtins.int,
Expand Down Expand Up @@ -1331,15 +1332,32 @@ class IggyClient:
self,
stream: builtins.str | builtins.int,
topic: builtins.str | builtins.int,
partitioning: builtins.int,
partitioning: Partitioning | builtins.int,
messages: list[SendMessage],
) -> collections.abc.Awaitable[SendMessagesResponse]:
r"""
Sends a list of messages to the specified topic.
Returns a SendMessagesResponse carrying the per-partition commit
confirmations, or a PyRuntimeError on failure. The confirmation list is
empty when the server reports no offsets, and the legacy server never
reports any.
Sends a batch of messages to a topic using the selected partitioning strategy.

Args:
stream: Stream identifier as `str | int`.
topic: Topic identifier as `str | int`.
partitioning: A `Partitioning` strategy or an integer partition ID.
Use `Partitioning.balanced()`, `Partitioning.partition_id(id)`, or
`Partitioning.messages_key(key)`. An integer is shorthand for
`Partitioning.partition_id(id)`.
messages: Messages to send as `list[SendMessage]`.

Returns:
An awaitable that resolves to `SendMessagesResponse`. Its confirmations
report the committed partition and batch base offset. The list is empty
when the server reports no offsets, including on the legacy server.

Raises:
ValueError: If a string stream or topic identifier is invalid.
TypeError: If `partitioning` or `messages` has an unsupported type.
OverflowError: If a numeric stream, topic, or partition ID is outside
the supported unsigned 32-bit range.
RuntimeError: If the request fails.
"""
def poll_messages(
self,
Expand Down Expand Up @@ -1639,6 +1657,37 @@ class Partition:
The number of messages in the partition.
"""

@typing.final
class Partitioning:
Comment thread
jiengup marked this conversation as resolved.
r"""
Defines how a batch of messages is assigned to a topic partition.
"""
@staticmethod
def balanced() -> Partitioning:
r"""
Routes the batch to one partition selected by round-robin.
"""
@staticmethod
def partition_id(partition_id: builtins.int) -> Partitioning:
r"""
Routes the batch to the specified partition.

`partition_id` must be between 0 and `2**32 - 1`. The topic must contain
that partition when the batch is sent.
"""
@staticmethod
def messages_key(key: builtins.str | bytes) -> Partitioning:
r"""
Routes the batch to one partition selected by hashing `key`.

`key` may be `str` or `bytes`. Strings are encoded as UTF-8; the encoded
key must contain between 1 and 255 bytes.

Raises:
ValueError: If the encoded key is empty or exceeds 255 bytes.
TypeError: If `key` is not `str` or `bytes`.
"""

@typing.final
class Permissions:
r"""
Expand Down Expand Up @@ -2103,8 +2152,8 @@ class Topic:
r"""
Options admission resolved for the keys the client did not send.

Same shape as `options`. These would have resolved differently under
another server configuration.
Same shape as [`Self::options`]. These would have resolved differently
under another server configuration.
"""

@typing.final
Expand Down Expand Up @@ -2168,8 +2217,8 @@ class TopicDetails:
r"""
Options admission resolved for the keys the client did not send.

Same shape as `options`. These would have resolved differently under
another server configuration.
Same shape as [`Self::options`]. These would have resolved differently
under another server configuration.
"""
@property
def partitions(self) -> builtins.list[Partition]:
Expand Down
35 changes: 27 additions & 8 deletions foreign/python/src/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@ use bytes::Bytes;
use iggy::prelude::{
AutoCommit as RustAutoCommit, Consumer as RustConsumer, IggyClient as RustIggyClient,
IggyExpiry as RustIggyExpiry, IggyMessage as RustMessage, MaxTopicSize as RustMaxTopicSize,
PollingStrategy as RustPollingStrategy, *,
Partitioning as RustPartitioning, PollingStrategy as RustPollingStrategy, *,
};
use pyo3::PyRef;
use pyo3::prelude::*;
Expand All @@ -39,6 +39,7 @@ use crate::consumer::{
use crate::duration::{py_delta_to_iggy_duration, reject_zero};
use crate::identifier::PyIdentifier;
use crate::options::OptionSpec as PyOptionSpec;
use crate::partitioning::PyPartitioning;
use crate::permissions::Permissions as PyPermissions;
use crate::receive_message::{PollingStrategy, ReceiveMessage};
use crate::send_message::{SendMessage, SendMessagesResponse as PySendMessagesResponse};
Expand Down Expand Up @@ -995,18 +996,36 @@ impl IggyClient {
})
}

/// Sends a list of messages to the specified topic.
/// Returns a SendMessagesResponse carrying the per-partition commit
/// confirmations, or a PyRuntimeError on failure. The confirmation list is
/// empty when the server reports no offsets, and the legacy server never
/// reports any.
/// Sends a batch of messages to a topic using the selected partitioning strategy.
///
/// Args:
/// stream: Stream identifier as `str | int`.
/// topic: Topic identifier as `str | int`.
/// partitioning: A `Partitioning` strategy or an integer partition ID.
/// Use `Partitioning.balanced()`, `Partitioning.partition_id(id)`, or
/// `Partitioning.messages_key(key)`. An integer is shorthand for
/// `Partitioning.partition_id(id)`.
/// messages: Messages to send as `list[SendMessage]`.
///
/// Returns:
/// An awaitable that resolves to `SendMessagesResponse`. Its confirmations
/// report the committed partition and batch base offset. The list is empty
/// when the server reports no offsets, including on the legacy server.
///
/// Raises:
/// ValueError: If a string stream or topic identifier is invalid.
/// TypeError: If `partitioning` or `messages` has an unsupported type.
/// OverflowError: If a numeric stream, topic, or partition ID is outside
/// the supported unsigned 32-bit range.
/// RuntimeError: If the request fails.
#[gen_stub(override_return_type(type_repr="collections.abc.Awaitable[SendMessagesResponse]", imports=("collections.abc")))]
fn send_messages<'a>(
&self,
py: Python<'a>,
stream: PyIdentifier,
topic: PyIdentifier,
partitioning: u32,
#[gen_stub(override_type(type_repr = "Partitioning | builtins.int"))]
partitioning: PyPartitioning,
#[gen_stub(override_type(type_repr = "list[SendMessage]"))] messages: &Bound<'_, PyList>,
) -> PyResult<Bound<'a, PyAny>> {
let messages: Vec<SendMessage> = messages
Expand All @@ -1023,7 +1042,7 @@ impl IggyClient {

let stream = Identifier::try_from(stream)?;
let topic = Identifier::try_from(topic)?;
let partitioning = Partitioning::partition_id(partitioning);
let partitioning = RustPartitioning::from(partitioning);
let inner = self.inner.clone();

future_into_py(py, async move {
Expand Down
3 changes: 3 additions & 0 deletions foreign/python/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ mod consumer;
mod duration;
mod identifier;
mod options;
mod partitioning;
mod permissions;
mod receive_message;
mod send_message;
Expand All @@ -36,6 +37,7 @@ use consumer::{
ConsumerGroupMember, IggyConsumer, ReceiveMessageIterator,
};
use options::OptionSpec;
use partitioning::Partitioning;
use permissions::{GlobalPermissions, Permissions, StreamPermissions, TopicPermissions};
use pyo3::prelude::*;
use receive_message::{PollingStrategy, ReceiveMessage};
Expand All @@ -62,6 +64,7 @@ fn apache_iggy(_py: Python, m: &Bound<'_, PyModule>) -> PyResult<()> {
m.add_class::<IggyExpiry>()?;
m.add_class::<MaxTopicSize>()?;
m.add_class::<OptionSpec>()?;
m.add_class::<Partitioning>()?;
m.add_class::<Partition>()?;
m.add_class::<Consumer>()?;
m.add_class::<ConsumerGroup>()?;
Expand Down
100 changes: 100 additions & 0 deletions foreign/python/src/partitioning.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,100 @@
// Licensed to the Apache Software Foundation (ASF) under one
// or more contributor license agreements. See the NOTICE file
// distributed with this work for additional information
// regarding copyright ownership. The ASF licenses this file
// to you under the Apache License, Version 2.0 (the
// "License"); you may not use this file except in compliance
// with the License. You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing,
// software distributed under the License is distributed on an
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
// KIND, either express or implied. See the License for the
// specific language governing permissions and limitations
// under the License.

use iggy::prelude::Partitioning as RustPartitioning;
use pyo3::{exceptions::PyValueError, prelude::*, types::PyBytes};
use pyo3_stub_gen::{
derive::{gen_stub_pyclass, gen_stub_pymethods},
impl_stub_type,
};

/// Defines how a batch of messages is assigned to a topic partition.
#[derive(Clone)]
#[pyclass(from_py_object)]
#[gen_stub_pyclass]
pub struct Partitioning {
pub(crate) inner: RustPartitioning,
}

#[gen_stub_pymethods]
#[pymethods]
impl Partitioning {
/// Routes the batch to one partition selected by round-robin.
#[staticmethod]
pub fn balanced() -> Self {
Self {
inner: RustPartitioning::balanced(),
}
}

/// Routes the batch to the specified partition.
Comment thread
jiengup marked this conversation as resolved.
///
/// `partition_id` must be between 0 and `2**32 - 1`. The topic must contain
/// that partition when the batch is sent.
#[staticmethod]
pub fn partition_id(partition_id: u32) -> Self {
Self {
inner: RustPartitioning::partition_id(partition_id),
}
}

/// Routes the batch to one partition selected by hashing `key`.
///
/// `key` may be `str` or `bytes`. Strings are encoded as UTF-8; the encoded
/// key must contain between 1 and 255 bytes.
///
/// Raises:
/// ValueError: If the encoded key is empty or exceeds 255 bytes.
/// TypeError: If `key` is not `str` or `bytes`.
#[staticmethod]
pub fn messages_key(py: Python<'_>, key: PyMessagesKey) -> PyResult<Self> {
let key = match key {
PyMessagesKey::String(key) => key.into_bytes(),
PyMessagesKey::Bytes(key) => key.extract::<Vec<u8>>(py)?,
};
let inner = RustPartitioning::messages_key(&key)
.map_err(|error| PyValueError::new_err(error.to_string()))?;
Ok(Self { inner })
}
}

#[derive(FromPyObject)]
pub enum PyMessagesKey {
#[pyo3(transparent, annotation = "str")]
String(String),
#[pyo3(transparent, annotation = "bytes")]
Bytes(Py<PyBytes>),
}
impl_stub_type!(PyMessagesKey = String | PyBytes);

#[derive(FromPyObject)]
pub(crate) enum PyPartitioning {
#[pyo3(transparent, annotation = "Partitioning")]
Strategy(Partitioning),
#[pyo3(transparent, annotation = "int")]
PartitionId(u32),
}
impl_stub_type!(PyPartitioning = Partitioning | isize);

impl From<PyPartitioning> for RustPartitioning {
fn from(partitioning: PyPartitioning) -> Self {
match partitioning {
PyPartitioning::Strategy(partitioning) => partitioning.inner,
PyPartitioning::PartitionId(partition_id) => Self::partition_id(partition_id),
}
}
}
Loading
Loading