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
8 changes: 7 additions & 1 deletion examples/python/basic/consumer.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, PollingStrategy, ReceiveMessage
from apache_iggy import IggyClient, IggyError, PollingStrategy, ReceiveMessage
from loguru import logger

STREAM_NAME = "sample-stream"
Expand Down Expand Up @@ -88,6 +88,12 @@ async def consume_messages(client: IggyClient):
handle_message(message)
n_consumed_batches += 1
await asyncio.sleep(interval)
except IggyError as error:
logger.error(
"Iggy error while consuming messages: "
f"[{error.code}] {error.name}: {error.message}"
)
break
except Exception as error:
logger.exception(f"Exception occurred while consuming messages: {error}")
break
Expand Down
15 changes: 14 additions & 1 deletion 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, IggyError, StreamDetails, TopicDetails
from apache_iggy import SendMessage as Message
from loguru import logger

Expand Down Expand Up @@ -69,6 +69,10 @@ async def init_system(client: IggyClient):
else:
logger.warning(f"Stream {stream.name} already exists with ID {stream.id}")

except IggyError as error:
logger.error(
f"Error creating stream: [{error.code}] {error.name}: {error.message}"
)
except Exception as error:
logger.error(f"Error creating stream: {error}")
logger.exception(error)
Expand All @@ -86,6 +90,10 @@ async def init_system(client: IggyClient):
logger.info("Topic was created successfully.")
else:
logger.warning(f"Topic {topic.name} already exists with ID {topic.id}")
except IggyError as error:
logger.error(
f"Error creating topic: [{error.code}] {error.name}: {error.message}"
)
except Exception as error:
logger.error(f"Error creating topic {error}")
logger.exception(error)
Expand Down Expand Up @@ -124,6 +132,11 @@ async def produce_messages(client: IggyClient):
f"Successfully sent batch of {messages_per_batch} messages. "
f"Batch ID: {current_id // messages_per_batch}"
)
except IggyError as error:
logger.error(
"Iggy error while sending messages: "
f"[{error.code}] {error.name}: {error.message}"
)
except Exception as error:
logger.error(f"Exception type: {type(error).__name__}, message: {error}")
logger.exception(error)
Expand Down
12 changes: 11 additions & 1 deletion examples/python/getting-started/consumer.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@
import typing
import urllib.parse

from apache_iggy import IggyClient, PollingStrategy, ReceiveMessage
from apache_iggy import IggyClient, IggyError, PollingStrategy, ReceiveMessage
from loguru import logger

STREAM_NAME = "sample-stream"
Expand Down Expand Up @@ -122,6 +122,10 @@ async def main():
await client.connect()
logger.info("Connected.")
await consume_messages(client)
except IggyError as error:
logger.error(
f"Iggy error in main: [{error.code}] {error.name}: {error.message}"
)
except Exception as error:
logger.exception(f"Exception occurred in main function: {error}")

Expand Down Expand Up @@ -157,6 +161,12 @@ async def consume_messages(client: IggyClient):
handle_message(message)
n_consumed_batches += 1
await asyncio.sleep(interval)
except IggyError as error:
logger.error(
"Iggy error while consuming messages: "
f"[{error.code}] {error.name}: {error.message}"
)
break
except Exception as error:
logger.exception(f"Exception occurred while consuming messages: {error}")
break
Expand Down
15 changes: 14 additions & 1 deletion examples/python/getting-started/producer.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@
import typing
import urllib.parse

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

Expand Down Expand Up @@ -135,6 +135,10 @@ async def init_system(client: IggyClient):
else:
logger.warning(f"Stream {stream.name} already exists with ID {stream.id}")

except IggyError as error:
logger.error(
f"Error creating stream: [{error.code}] {error.name}: {error.message}"
)
except Exception as error:
logger.error(f"Error creating stream: {error}")
logger.exception(error)
Expand All @@ -152,6 +156,10 @@ async def init_system(client: IggyClient):
logger.info("Topic was created successfully.")
else:
logger.warning(f"Topic {topic.name} already exists with ID {topic.id}")
except IggyError as error:
logger.error(
f"Error creating topic: [{error.code}] {error.name}: {error.message}"
)
except Exception as error:
logger.error(f"Error creating topic {error}")
logger.exception(error)
Expand Down Expand Up @@ -190,6 +198,11 @@ async def produce_messages(client: IggyClient):
f"Successfully sent batch of {messages_per_batch} messages. "
f"Batch ID: {current_id // messages_per_batch}"
)
except IggyError as error:
logger.error(
"Iggy error while sending messages: "
f"[{error.code}] {error.name}: {error.message}"
)
except Exception as error:
logger.error(f"Exception type: {type(error).__name__}, message: {error}")
logger.exception(error)
Expand Down
3 changes: 2 additions & 1 deletion foreign/python/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -48,7 +48,8 @@ cargo run --bin iggy-server -- --with-default-root-credentials --fresh
# Using uv:
uv sync --all-extras
uv run maturin develop
uv run pytest tests/ -v # Run tests (requires iggy-server running)
uv run pytest tests/ -v # Run all tests (for release versions, requires iggy-server running)
uv run --no-sync pytest tests/ -v # Run tests without syncing (for development, implicit uv run's syncing overwrites the installations made by maturin develop)

# Using pip:
python3 -m venv .venv
Expand Down
19 changes: 19 additions & 0 deletions foreign/python/apache_iggy.pyi
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@ __all__ = [
"ConsumerGroupMember",
"IggyClient",
"IggyConsumer",
"IggyError",
"PollingStrategy",
"ReceiveMessage",
"SendMessage",
Expand Down Expand Up @@ -709,6 +710,24 @@ class IggyConsumer:
Returns an awaitable that completes when shutdown is signaled or a PyRuntimeError on failure.
"""

@typing.final
class IggyError(builtins.Exception):
r"""
A Python class representing the Rust's IggyError.
Allows transparent representation of Iggy-specific errors.
"""
@property
def code(self) -> builtins.int: ...
@property
def name(self) -> builtins.str: ...
@property
def message(self) -> builtins.str: ...
def __new__(
cls, code: builtins.int, name: builtins.str, message: builtins.str
) -> IggyError: ...
def __str__(self) -> builtins.str: ...
def __repr__(self) -> builtins.str: ...

class PollingStrategy:
@typing.final
class Offset(PollingStrategy):
Expand Down
Loading
Loading